Skip to main content

parquet/file/
serialized_reader.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 implementations of the reader traits FileReader, RowGroupReader and PageReader
19//! Also contains implementations of the ChunkReader for files (with buffering) and byte arrays (RAM)
20
21use crate::basic::{PageType, Type};
22use crate::bloom_filter::Sbbf;
23use crate::column::page::{Page, PageMetadata, PageReader};
24use crate::compression::{Codec, create_codec};
25#[cfg(feature = "encryption")]
26use crate::encryption::decrypt::{CryptoContext, read_and_decrypt};
27use crate::errors::{ParquetError, Result};
28use crate::file::metadata::thrift::PageHeader;
29use crate::file::page_index::offset_index::{OffsetIndexMetaData, PageLocation};
30use crate::file::statistics;
31use crate::file::{
32    metadata::*,
33    properties::{ReaderProperties, ReaderPropertiesPtr},
34    reader::*,
35};
36#[cfg(feature = "encryption")]
37use crate::parquet_thrift::ThriftSliceInputProtocol;
38use crate::parquet_thrift::{ReadThrift, ThriftReadInputProtocol};
39use crate::record::Row;
40use crate::record::reader::RowIter;
41use crate::schema::types::{SchemaDescPtr, Type as SchemaType};
42use bytes::Bytes;
43use std::collections::VecDeque;
44use std::{fs::File, io::Read, path::Path, sync::Arc};
45
46impl TryFrom<File> for SerializedFileReader<File> {
47    type Error = ParquetError;
48
49    fn try_from(file: File) -> Result<Self> {
50        Self::new(file)
51    }
52}
53
54impl TryFrom<&Path> for SerializedFileReader<File> {
55    type Error = ParquetError;
56
57    fn try_from(path: &Path) -> Result<Self> {
58        let file = File::open(path)?;
59        Self::try_from(file)
60    }
61}
62
63impl TryFrom<String> for SerializedFileReader<File> {
64    type Error = ParquetError;
65
66    fn try_from(path: String) -> Result<Self> {
67        Self::try_from(Path::new(&path))
68    }
69}
70
71impl TryFrom<&str> for SerializedFileReader<File> {
72    type Error = ParquetError;
73
74    fn try_from(path: &str) -> Result<Self> {
75        Self::try_from(Path::new(&path))
76    }
77}
78
79/// Conversion into a [`RowIter`]
80/// using the full file schema over all row groups.
81impl IntoIterator for SerializedFileReader<File> {
82    type Item = Result<Row>;
83    type IntoIter = RowIter<'static>;
84
85    fn into_iter(self) -> Self::IntoIter {
86        RowIter::from_file_into(Box::new(self))
87    }
88}
89
90// ----------------------------------------------------------------------
91// Implementations of file & row group readers
92
93/// A serialized implementation for Parquet [`FileReader`].
94pub struct SerializedFileReader<R: ChunkReader> {
95    chunk_reader: Arc<R>,
96    metadata: Arc<ParquetMetaData>,
97    props: ReaderPropertiesPtr,
98}
99
100/// A predicate for filtering row groups, invoked with the metadata and index
101/// of each row group in the file. Only row groups for which the predicate
102/// evaluates to `true` will be scanned
103pub type ReadGroupPredicate = Box<dyn FnMut(&RowGroupMetaData, usize) -> bool>;
104
105/// A builder for [`ReadOptions`].
106/// For the predicates that are added to the builder,
107/// they will be chained using 'AND' to filter the row groups.
108#[derive(Default)]
109pub struct ReadOptionsBuilder {
110    predicates: Vec<ReadGroupPredicate>,
111    enable_page_index: bool,
112    props: Option<ReaderProperties>,
113    metadata_options: ParquetMetaDataOptions,
114}
115
116impl ReadOptionsBuilder {
117    /// New builder
118    pub fn new() -> Self {
119        Self::default()
120    }
121
122    /// Add a predicate on row group metadata to the reading option,
123    /// Filter only row groups that match the predicate criteria
124    pub fn with_predicate(mut self, predicate: ReadGroupPredicate) -> Self {
125        self.predicates.push(predicate);
126        self
127    }
128
129    /// Add a range predicate on filtering row groups if their midpoints are within
130    /// the Closed-Open range `[start..end) {x | start <= x < end}`
131    pub fn with_range(mut self, start: i64, end: i64) -> Self {
132        assert!(start < end);
133        let predicate = move |rg: &RowGroupMetaData, _: usize| {
134            let mid = get_midpoint_offset(rg);
135            mid >= start && mid < end
136        };
137        self.predicates.push(Box::new(predicate));
138        self
139    }
140
141    /// Enable reading the page index structures described in
142    /// "[Column Index] Layout to Support Page Skipping"
143    ///
144    /// [Column Index]: https://github.com/apache/parquet-format/blob/master/PageIndex.md
145    pub fn with_page_index(mut self) -> Self {
146        self.enable_page_index = true;
147        self
148    }
149
150    /// Set the [`ReaderProperties`] configuration.
151    pub fn with_reader_properties(mut self, properties: ReaderProperties) -> Self {
152        self.props = Some(properties);
153        self
154    }
155
156    /// Provide a Parquet schema to use when decoding the metadata. The schema in the Parquet
157    /// footer will be skipped.
158    pub fn with_parquet_schema(mut self, schema: SchemaDescPtr) -> Self {
159        self.metadata_options.set_schema(schema);
160        self
161    }
162
163    /// Set whether to convert the [`encoding_stats`] in the Parquet `ColumnMetaData` to a bitmask
164    /// (defaults to `false`).
165    ///
166    /// See [`ColumnChunkMetaData::page_encoding_stats_mask`] for an explanation of why this
167    /// might be desirable.
168    ///
169    /// [`encoding_stats`]:
170    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
171    pub fn with_encoding_stats_as_mask(mut self, val: bool) -> Self {
172        self.metadata_options.set_encoding_stats_as_mask(val);
173        self
174    }
175
176    /// Sets the decoding policy for [`encoding_stats`] in the Parquet `ColumnMetaData`.
177    ///
178    /// [`encoding_stats`]:
179    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L917
180    pub fn with_encoding_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
181        self.metadata_options.set_encoding_stats_policy(policy);
182        self
183    }
184
185    /// Sets the decoding policy for [`statistics`] in the Parquet `ColumnMetaData`.
186    ///
187    /// [`statistics`]:
188    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L912
189    pub fn with_column_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
190        self.metadata_options.set_column_stats_policy(policy);
191        self
192    }
193
194    /// Sets the decoding policy for [`size_statistics`] in the Parquet `ColumnMetaData`.
195    ///
196    /// [`size_statistics`]:
197    /// https://github.com/apache/parquet-format/blob/786142e26740487930ddc3ec5e39d780bd930907/src/main/thrift/parquet.thrift#L936
198    pub fn with_size_stats_policy(mut self, policy: ParquetStatisticsPolicy) -> Self {
199        self.metadata_options.set_size_stats_policy(policy);
200        self
201    }
202
203    /// Seal the builder and return the read options
204    pub fn build(self) -> ReadOptions {
205        let props = self
206            .props
207            .unwrap_or_else(|| ReaderProperties::builder().build());
208        ReadOptions {
209            predicates: self.predicates,
210            enable_page_index: self.enable_page_index,
211            props,
212            metadata_options: self.metadata_options,
213        }
214    }
215}
216
217/// A collection of options for reading a Parquet file.
218///
219/// Predicates are currently only supported on row group metadata.
220/// All predicates will be chained using 'AND' to filter the row groups.
221pub struct ReadOptions {
222    predicates: Vec<ReadGroupPredicate>,
223    enable_page_index: bool,
224    props: ReaderProperties,
225    metadata_options: ParquetMetaDataOptions,
226}
227
228impl<R: 'static + ChunkReader> SerializedFileReader<R> {
229    /// Creates file reader from a Parquet file.
230    /// Returns an error if the Parquet file does not exist or is corrupt.
231    pub fn new(chunk_reader: R) -> Result<Self> {
232        let metadata = ParquetMetaDataReader::new().parse_and_finish(&chunk_reader)?;
233        let props = Arc::new(ReaderProperties::builder().build());
234        Ok(Self {
235            chunk_reader: Arc::new(chunk_reader),
236            metadata: Arc::new(metadata),
237            props,
238        })
239    }
240
241    /// Creates file reader from a Parquet file with read options.
242    /// Returns an error if the Parquet file does not exist or is corrupt.
243    pub fn new_with_options(chunk_reader: R, options: ReadOptions) -> Result<Self> {
244        let mut metadata_builder = ParquetMetaDataReader::new()
245            .with_metadata_options(Some(options.metadata_options.clone()))
246            .parse_and_finish(&chunk_reader)?
247            .into_builder();
248        let mut predicates = options.predicates;
249
250        // Filter row groups based on the predicates
251        for (i, rg_meta) in metadata_builder.take_row_groups().into_iter().enumerate() {
252            let mut keep = true;
253            for predicate in &mut predicates {
254                if !predicate(&rg_meta, i) {
255                    keep = false;
256                    break;
257                }
258            }
259            if keep {
260                metadata_builder = metadata_builder.add_row_group(rg_meta);
261            }
262        }
263
264        let mut metadata = metadata_builder.build();
265
266        // If page indexes are desired, build them with the filtered set of row groups
267        if options.enable_page_index {
268            let mut reader = ParquetMetaDataReader::new_with_metadata(metadata)
269                .with_page_index_policy(PageIndexPolicy::Required);
270            reader.read_page_indexes(&chunk_reader)?;
271            metadata = reader.finish()?;
272        }
273
274        Ok(Self {
275            chunk_reader: Arc::new(chunk_reader),
276            metadata: Arc::new(metadata),
277            props: Arc::new(options.props),
278        })
279    }
280}
281
282/// Get midpoint offset for a row group
283fn get_midpoint_offset(meta: &RowGroupMetaData) -> i64 {
284    let col = meta.column(0);
285    let mut offset = col.data_page_offset();
286    if let Some(dic_offset) = col.dictionary_page_offset()
287        && offset > dic_offset
288    {
289        offset = dic_offset
290    };
291    offset + meta.compressed_size() / 2
292}
293
294impl<R: 'static + ChunkReader> FileReader for SerializedFileReader<R> {
295    fn metadata(&self) -> &ParquetMetaData {
296        &self.metadata
297    }
298
299    fn num_row_groups(&self) -> usize {
300        self.metadata.num_row_groups()
301    }
302
303    fn get_row_group(&self, i: usize) -> Result<Box<dyn RowGroupReader + '_>> {
304        let row_group_metadata = self.metadata.row_group(i);
305        // Row groups should be processed sequentially.
306        let props = Arc::clone(&self.props);
307        let f = Arc::clone(&self.chunk_reader);
308        Ok(Box::new(SerializedRowGroupReader::new(
309            f,
310            row_group_metadata,
311            self.metadata.offset_index().map(|x| x[i].as_slice()),
312            props,
313        )?))
314    }
315
316    fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
317        RowIter::from_file(projection, self)
318    }
319}
320
321/// A serialized implementation for Parquet [`RowGroupReader`].
322pub struct SerializedRowGroupReader<'a, R: ChunkReader> {
323    chunk_reader: Arc<R>,
324    metadata: &'a RowGroupMetaData,
325    offset_index: Option<&'a [OffsetIndexMetaData]>,
326    props: ReaderPropertiesPtr,
327    bloom_filters: Vec<Option<Sbbf>>,
328}
329
330impl<'a, R: ChunkReader> SerializedRowGroupReader<'a, R> {
331    /// Creates new row group reader from a file, row group metadata and custom config.
332    pub fn new(
333        chunk_reader: Arc<R>,
334        metadata: &'a RowGroupMetaData,
335        offset_index: Option<&'a [OffsetIndexMetaData]>,
336        props: ReaderPropertiesPtr,
337    ) -> Result<Self> {
338        let bloom_filters = if props.read_bloom_filter() {
339            metadata
340                .columns()
341                .iter()
342                .map(|col| Sbbf::read_from_column_chunk(col, &*chunk_reader))
343                .collect::<Result<Vec<_>>>()?
344        } else {
345            std::iter::repeat_n(None, metadata.columns().len()).collect()
346        };
347        Ok(Self {
348            chunk_reader,
349            metadata,
350            offset_index,
351            props,
352            bloom_filters,
353        })
354    }
355}
356
357impl<R: 'static + ChunkReader> RowGroupReader for SerializedRowGroupReader<'_, R> {
358    fn metadata(&self) -> &RowGroupMetaData {
359        self.metadata
360    }
361
362    fn num_columns(&self) -> usize {
363        self.metadata.num_columns()
364    }
365
366    // TODO: fix PARQUET-816
367    fn get_column_page_reader(&self, i: usize) -> Result<Box<dyn PageReader>> {
368        let col = self.metadata.column(i);
369
370        let page_locations = self.offset_index.map(|x| x[i].page_locations.clone());
371
372        let props = Arc::clone(&self.props);
373        Ok(Box::new(SerializedPageReader::new_with_properties(
374            Arc::clone(&self.chunk_reader),
375            col,
376            usize::try_from(self.metadata.num_rows())?,
377            page_locations,
378            props,
379        )?))
380    }
381
382    /// get bloom filter for the `i`th column
383    fn get_column_bloom_filter(&self, i: usize) -> Option<&Sbbf> {
384        self.bloom_filters[i].as_ref()
385    }
386
387    fn get_row_iter(&self, projection: Option<SchemaType>) -> Result<RowIter<'_>> {
388        RowIter::from_row_group(projection, self)
389    }
390}
391
392/// Decodes a [`Page`] from the provided `buffer`
393pub(crate) fn decode_page(
394    page_header: PageHeader,
395    buffer: Bytes,
396    physical_type: Type,
397    decompressor: Option<&mut Box<dyn Codec>>,
398) -> Result<Page> {
399    // Verify the 32-bit CRC checksum of the page
400    #[cfg(feature = "crc")]
401    if let Some(expected_crc) = page_header.crc {
402        let crc = crc32fast::hash(&buffer);
403        if crc != expected_crc as u32 {
404            return Err(general_err!("Page CRC checksum mismatch"));
405        }
406    }
407
408    // When processing data page v2, depending on enabled compression for the
409    // page, we should account for uncompressed data ('offset') of
410    // repetition and definition levels.
411    //
412    // We always use 0 offset for other pages other than v2, `true` flag means
413    // that compression will be applied if decompressor is defined
414    let (offset, can_decompress): (usize, bool) = match page_header.data_page_header_v2 {
415        Some(ref header_v2) => {
416            if header_v2.definition_levels_byte_length < 0
417                || header_v2.repetition_levels_byte_length < 0
418                || header_v2.definition_levels_byte_length + header_v2.repetition_levels_byte_length
419                    > page_header.uncompressed_page_size
420            {
421                return Err(general_err!(
422                    "DataPage v2 header contains implausible values \
423                        for definition_levels_byte_length ({}) \
424                        and repetition_levels_byte_length ({}) \
425                        given DataPage header provides uncompressed_page_size ({})",
426                    header_v2.definition_levels_byte_length,
427                    header_v2.repetition_levels_byte_length,
428                    page_header.uncompressed_page_size
429                ));
430            }
431            (
432                usize::try_from(
433                    header_v2.definition_levels_byte_length
434                        + header_v2.repetition_levels_byte_length,
435                )?,
436                // When is_compressed flag is missing the page is considered compressed
437                header_v2.is_compressed.unwrap_or(true),
438            )
439        }
440        None => (0, true),
441    };
442
443    let buffer = match decompressor {
444        Some(decompressor) if can_decompress => {
445            let uncompressed_page_size = usize::try_from(page_header.uncompressed_page_size)?;
446            if offset > buffer.len() || offset > uncompressed_page_size {
447                return Err(general_err!("Invalid page header"));
448            }
449            let decompressed_size = uncompressed_page_size - offset;
450            let mut decompressed = Vec::with_capacity(uncompressed_page_size);
451            decompressed.extend_from_slice(&buffer[..offset]);
452            // decompressed size of zero corresponds to a page with no non-null values
453            // see https://github.com/apache/parquet-format/blob/master/README.md#data-pages
454            if decompressed_size > 0 {
455                let compressed = &buffer[offset..];
456                decompressor.decompress(compressed, &mut decompressed, Some(decompressed_size))?;
457            }
458
459            if decompressed.len() != uncompressed_page_size {
460                return Err(general_err!(
461                    "Actual decompressed size doesn't match the expected one ({} vs {})",
462                    decompressed.len(),
463                    uncompressed_page_size
464                ));
465            }
466
467            Bytes::from(decompressed)
468        }
469        _ => buffer,
470    };
471
472    let result = match page_header.r#type {
473        PageType::DICTIONARY_PAGE => {
474            let dict_header = page_header.dictionary_page_header.as_ref().ok_or_else(|| {
475                ParquetError::General("Missing dictionary page header".to_string())
476            })?;
477            let is_sorted = dict_header.is_sorted.unwrap_or(false);
478            Page::DictionaryPage {
479                buf: buffer,
480                num_values: dict_header.num_values.try_into()?,
481                encoding: dict_header.encoding,
482                is_sorted,
483            }
484        }
485        PageType::DATA_PAGE => {
486            let header = page_header
487                .data_page_header
488                .ok_or_else(|| ParquetError::General("Missing V1 data page header".to_string()))?;
489            Page::DataPage {
490                buf: buffer,
491                num_values: header.num_values.try_into()?,
492                encoding: header.encoding,
493                def_level_encoding: header.definition_level_encoding,
494                rep_level_encoding: header.repetition_level_encoding,
495                statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
496            }
497        }
498        PageType::DATA_PAGE_V2 => {
499            let header = page_header
500                .data_page_header_v2
501                .ok_or_else(|| ParquetError::General("Missing V2 data page header".to_string()))?;
502            let is_compressed = header.is_compressed.unwrap_or(true);
503            Page::DataPageV2 {
504                buf: buffer,
505                num_values: header.num_values.try_into()?,
506                encoding: header.encoding,
507                num_nulls: header.num_nulls.try_into()?,
508                num_rows: header.num_rows.try_into()?,
509                def_levels_byte_len: header.definition_levels_byte_length.try_into()?,
510                rep_levels_byte_len: header.repetition_levels_byte_length.try_into()?,
511                is_compressed,
512                statistics: statistics::from_thrift_page_stats(physical_type, header.statistics)?,
513            }
514        }
515        _ => {
516            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
517            return Err(general_err!(
518                "Page type {:?} is not supported",
519                page_header.r#type
520            ));
521        }
522    };
523
524    Ok(result)
525}
526
527enum SerializedPageReaderState {
528    Values {
529        /// The current byte offset in the reader
530        /// Note that offset is u64 (i.e., not usize) to support 32-bit architectures such as WASM
531        offset: u64,
532
533        /// The length of the chunk in bytes
534        /// Note that remaining_bytes is u64 (i.e., not usize) to support 32-bit architectures such as WASM
535        remaining_bytes: u64,
536
537        // If the next page header has already been "peeked", we will cache it and it`s length here
538        next_page_header: Option<Box<PageHeader>>,
539
540        /// The index of the data page within this column chunk
541        page_index: usize,
542
543        /// Whether the next page is expected to be a dictionary page
544        require_dictionary: bool,
545    },
546    Pages {
547        /// Remaining page locations
548        page_locations: VecDeque<PageLocation>,
549        /// Remaining dictionary location if any
550        dictionary_page: Option<PageLocation>,
551        /// The total number of rows in this column chunk
552        total_rows: usize,
553        /// The index of the data page within this column chunk
554        page_index: usize,
555    },
556}
557
558#[derive(Default)]
559struct SerializedPageReaderContext {
560    /// Controls decoding of page-level statistics
561    read_stats: bool,
562    /// Crypto context carrying objects required for decryption
563    #[cfg(feature = "encryption")]
564    crypto_context: Option<Arc<CryptoContext>>,
565}
566
567/// A serialized implementation for Parquet [`PageReader`].
568pub struct SerializedPageReader<R: ChunkReader> {
569    /// The chunk reader
570    reader: Arc<R>,
571
572    /// The compression codec for this column chunk. Only set for non-PLAIN codec.
573    decompressor: Option<Box<dyn Codec>>,
574
575    /// Column chunk type.
576    physical_type: Type,
577
578    state: SerializedPageReaderState,
579
580    context: SerializedPageReaderContext,
581}
582
583impl<R: ChunkReader> SerializedPageReader<R> {
584    /// Creates a new serialized page reader from a chunk reader and metadata
585    pub fn new(
586        reader: Arc<R>,
587        column_chunk_metadata: &ColumnChunkMetaData,
588        total_rows: usize,
589        page_locations: Option<Vec<PageLocation>>,
590    ) -> Result<Self> {
591        let props = Arc::new(ReaderProperties::builder().build());
592        SerializedPageReader::new_with_properties(
593            reader,
594            column_chunk_metadata,
595            total_rows,
596            page_locations,
597            props,
598        )
599    }
600
601    /// Stub No-op implementation when encryption is disabled.
602    #[cfg(all(feature = "arrow", not(feature = "encryption")))]
603    pub(crate) fn add_crypto_context(
604        self,
605        _rg_idx: usize,
606        _column_idx: usize,
607        _parquet_meta_data: &ParquetMetaData,
608        _column_chunk_metadata: &ColumnChunkMetaData,
609    ) -> Result<SerializedPageReader<R>> {
610        Ok(self)
611    }
612
613    /// Adds any necessary crypto context to this page reader, if encryption is enabled.
614    #[cfg(feature = "encryption")]
615    pub(crate) fn add_crypto_context(
616        mut self,
617        rg_idx: usize,
618        column_idx: usize,
619        parquet_meta_data: &ParquetMetaData,
620        column_chunk_metadata: &ColumnChunkMetaData,
621    ) -> Result<SerializedPageReader<R>> {
622        let Some(file_decryptor) = parquet_meta_data.file_decryptor() else {
623            return Ok(self);
624        };
625        let Some(crypto_metadata) = column_chunk_metadata.crypto_metadata() else {
626            return Ok(self);
627        };
628        let crypto_context =
629            CryptoContext::for_column(file_decryptor, crypto_metadata, rg_idx, column_idx)?;
630        self.context.crypto_context = Some(Arc::new(crypto_context));
631        Ok(self)
632    }
633
634    /// Creates a new serialized page with custom options.
635    pub fn new_with_properties(
636        reader: Arc<R>,
637        meta: &ColumnChunkMetaData,
638        total_rows: usize,
639        page_locations: Option<Vec<PageLocation>>,
640        props: ReaderPropertiesPtr,
641    ) -> Result<Self> {
642        let decompressor = create_codec(meta.compression(), props.codec_options())?;
643        let (start, len) = meta.byte_range();
644
645        let state = match page_locations {
646            Some(locations) => {
647                // If the offset of the first page doesn't match the start of the column chunk
648                // then the preceding space must contain a dictionary page.
649                let dictionary_page = match locations.first() {
650                    Some(dict_offset) if dict_offset.offset as u64 != start => Some(PageLocation {
651                        offset: start as i64,
652                        compressed_page_size: (dict_offset.offset as u64 - start) as i32,
653                        first_row_index: 0,
654                    }),
655                    _ => None,
656                };
657
658                SerializedPageReaderState::Pages {
659                    page_locations: locations.into(),
660                    dictionary_page,
661                    total_rows,
662                    page_index: 0,
663                }
664            }
665            None => SerializedPageReaderState::Values {
666                offset: start,
667                remaining_bytes: len,
668                next_page_header: None,
669                page_index: 0,
670                require_dictionary: meta.dictionary_page_offset().is_some(),
671            },
672        };
673        let mut context = SerializedPageReaderContext::default();
674        if props.read_page_stats() {
675            context.read_stats = true;
676        }
677        Ok(Self {
678            reader,
679            decompressor,
680            state,
681            physical_type: meta.column_type(),
682            context,
683        })
684    }
685
686    /// Similar to `peek_next_page`, but returns the offset of the next page instead of the page metadata.
687    /// Unlike page metadata, an offset can uniquely identify a page.
688    ///
689    /// This is used when we need to read parquet with row-filter, and we don't want to decompress the page twice.
690    /// This function allows us to check if the next page is being cached or read previously.
691    #[cfg(test)]
692    fn peek_next_page_offset(&mut self) -> Result<Option<u64>> {
693        match &mut self.state {
694            SerializedPageReaderState::Values {
695                offset,
696                remaining_bytes,
697                next_page_header,
698                page_index,
699                require_dictionary,
700            } => {
701                loop {
702                    if *remaining_bytes == 0 {
703                        return Ok(None);
704                    }
705                    return if let Some(header) = next_page_header.as_ref() {
706                        if let Ok(_page_meta) = PageMetadata::try_from(&**header) {
707                            Ok(Some(*offset))
708                        } else {
709                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
710                            *next_page_header = None;
711                            continue;
712                        }
713                    } else {
714                        let mut read = self.reader.get_read(*offset)?;
715                        let (header_len, header) = Self::read_page_header_len(
716                            &self.context,
717                            &mut read,
718                            *page_index,
719                            *require_dictionary,
720                        )?;
721                        *offset += header_len as u64;
722                        *remaining_bytes -= header_len as u64;
723                        let page_meta = if let Ok(_page_meta) = PageMetadata::try_from(&header) {
724                            Ok(Some(*offset))
725                        } else {
726                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
727                            continue;
728                        };
729                        *next_page_header = Some(Box::new(header));
730                        page_meta
731                    };
732                }
733            }
734            SerializedPageReaderState::Pages {
735                page_locations,
736                dictionary_page,
737                ..
738            } => {
739                if let Some(page) = dictionary_page {
740                    Ok(Some(page.offset as u64))
741                } else if let Some(page) = page_locations.front() {
742                    Ok(Some(page.offset as u64))
743                } else {
744                    Ok(None)
745                }
746            }
747        }
748    }
749
750    fn read_page_header_len<T: Read>(
751        context: &SerializedPageReaderContext,
752        input: &mut T,
753        page_index: usize,
754        dictionary_page: bool,
755    ) -> Result<(usize, PageHeader)> {
756        /// A wrapper around a [`std::io::Read`] that keeps track of the bytes read
757        struct TrackedRead<R> {
758            inner: R,
759            bytes_read: usize,
760        }
761
762        impl<R: Read> Read for TrackedRead<R> {
763            fn read(&mut self, buf: &mut [u8]) -> std::io::Result<usize> {
764                let v = self.inner.read(buf)?;
765                self.bytes_read += v;
766                Ok(v)
767            }
768        }
769
770        let mut tracked = TrackedRead {
771            inner: input,
772            bytes_read: 0,
773        };
774        let header = context.read_page_header(&mut tracked, page_index, dictionary_page)?;
775        Ok((tracked.bytes_read, header))
776    }
777
778    fn read_page_header_len_from_bytes(
779        context: &SerializedPageReaderContext,
780        buffer: &[u8],
781        page_index: usize,
782        dictionary_page: bool,
783    ) -> Result<(usize, PageHeader)> {
784        let mut input = std::io::Cursor::new(buffer);
785        let header = context.read_page_header(&mut input, page_index, dictionary_page)?;
786        let header_len = input.position() as usize;
787        Ok((header_len, header))
788    }
789}
790
791#[cfg(not(feature = "encryption"))]
792impl SerializedPageReaderContext {
793    fn read_page_header<T: Read>(
794        &self,
795        input: &mut T,
796        _page_index: usize,
797        _dictionary_page: bool,
798    ) -> Result<PageHeader> {
799        let mut prot = ThriftReadInputProtocol::new(input);
800        if self.read_stats {
801            Ok(PageHeader::read_thrift(&mut prot)?)
802        } else {
803            Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
804        }
805    }
806
807    fn decrypt_page_data<T>(
808        &self,
809        buffer: T,
810        _page_index: usize,
811        _dictionary_page: bool,
812    ) -> Result<T> {
813        Ok(buffer)
814    }
815}
816
817#[cfg(feature = "encryption")]
818impl SerializedPageReaderContext {
819    fn read_page_header<T: Read>(
820        &self,
821        input: &mut T,
822        page_index: usize,
823        dictionary_page: bool,
824    ) -> Result<PageHeader> {
825        match self.page_crypto_context(page_index, dictionary_page) {
826            None => {
827                let mut prot = ThriftReadInputProtocol::new(input);
828                if self.read_stats {
829                    Ok(PageHeader::read_thrift(&mut prot)?)
830                } else {
831                    use crate::file::metadata::thrift::PageHeader;
832
833                    Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
834                }
835            }
836            Some(page_crypto_context) => {
837                let data_decryptor = page_crypto_context.data_decryptor();
838                let aad = page_crypto_context.create_page_header_aad()?;
839
840                let buf = read_and_decrypt(data_decryptor, input, aad.as_ref()).map_err(|_| {
841                    ParquetError::General(format!(
842                        "Error decrypting page header for column {}, decryption key may be wrong",
843                        page_crypto_context.column_ordinal
844                    ))
845                })?;
846
847                let mut prot = ThriftSliceInputProtocol::new(buf.as_slice());
848                if self.read_stats {
849                    Ok(PageHeader::read_thrift(&mut prot)?)
850                } else {
851                    Ok(PageHeader::read_thrift_without_stats(&mut prot)?)
852                }
853            }
854        }
855    }
856
857    fn decrypt_page_data<T>(&self, buffer: T, page_index: usize, dictionary_page: bool) -> Result<T>
858    where
859        T: AsRef<[u8]>,
860        T: From<Vec<u8>>,
861    {
862        let page_crypto_context = self.page_crypto_context(page_index, dictionary_page);
863        if let Some(page_crypto_context) = page_crypto_context {
864            let decryptor = page_crypto_context.data_decryptor();
865            let aad = page_crypto_context.create_page_aad()?;
866            let decrypted = decryptor.decrypt(buffer.as_ref(), &aad)?;
867            Ok(T::from(decrypted))
868        } else {
869            Ok(buffer)
870        }
871    }
872
873    fn page_crypto_context(
874        &self,
875        page_index: usize,
876        dictionary_page: bool,
877    ) -> Option<Arc<CryptoContext>> {
878        self.crypto_context.as_ref().map(|c| {
879            Arc::new(if dictionary_page {
880                c.for_dictionary_page()
881            } else {
882                c.with_page_ordinal(page_index)
883            })
884        })
885    }
886}
887
888impl<R: ChunkReader> Iterator for SerializedPageReader<R> {
889    type Item = Result<Page>;
890
891    fn next(&mut self) -> Option<Self::Item> {
892        self.get_next_page().transpose()
893    }
894}
895
896fn verify_page_header_len(header_len: usize, remaining_bytes: u64) -> Result<()> {
897    if header_len as u64 > remaining_bytes {
898        return Err(eof_err!("Invalid page header"));
899    }
900    Ok(())
901}
902
903fn verify_page_size(
904    compressed_size: i32,
905    uncompressed_size: i32,
906    remaining_bytes: u64,
907) -> Result<()> {
908    // The page's compressed size should not exceed the remaining bytes that are
909    // available to read. The page's uncompressed size is the expected size
910    // after decompression, which can never be negative.
911    if compressed_size < 0 || compressed_size as u64 > remaining_bytes || uncompressed_size < 0 {
912        return Err(eof_err!("Invalid page header"));
913    }
914    Ok(())
915}
916
917impl<R: ChunkReader> PageReader for SerializedPageReader<R> {
918    fn get_next_page(&mut self) -> Result<Option<Page>> {
919        loop {
920            let page = match &mut self.state {
921                SerializedPageReaderState::Values {
922                    offset,
923                    remaining_bytes: remaining,
924                    next_page_header,
925                    page_index,
926                    require_dictionary,
927                } => {
928                    if *remaining == 0 {
929                        return Ok(None);
930                    }
931
932                    let mut read = self.reader.get_read(*offset)?;
933                    let header = if let Some(header) = next_page_header.take() {
934                        *header
935                    } else {
936                        let (header_len, header) = Self::read_page_header_len(
937                            &self.context,
938                            &mut read,
939                            *page_index,
940                            *require_dictionary,
941                        )?;
942                        verify_page_header_len(header_len, *remaining)?;
943                        *offset += header_len as u64;
944                        *remaining -= header_len as u64;
945                        header
946                    };
947                    verify_page_size(
948                        header.compressed_page_size,
949                        header.uncompressed_page_size,
950                        *remaining,
951                    )?;
952                    let data_len = header.compressed_page_size as usize;
953                    let data_start = *offset;
954                    *offset += data_len as u64;
955                    *remaining -= data_len as u64;
956
957                    if header.r#type == PageType::INDEX_PAGE {
958                        continue;
959                    }
960
961                    let buffer = self.reader.get_bytes(data_start, data_len)?;
962
963                    let buffer =
964                        self.context
965                            .decrypt_page_data(buffer, *page_index, *require_dictionary)?;
966
967                    let page = decode_page(
968                        header,
969                        buffer,
970                        self.physical_type,
971                        self.decompressor.as_mut(),
972                    )?;
973                    if page.is_data_page() {
974                        *page_index += 1;
975                    } else if page.is_dictionary_page() {
976                        *require_dictionary = false;
977                    }
978                    page
979                }
980                SerializedPageReaderState::Pages {
981                    page_locations,
982                    dictionary_page,
983                    page_index,
984                    ..
985                } => {
986                    let (front, is_dictionary_page) = match dictionary_page.take() {
987                        Some(front) => (front, true),
988                        None => match page_locations.pop_front() {
989                            Some(front) => (front, false),
990                            None => return Ok(None),
991                        },
992                    };
993
994                    let page_len = usize::try_from(front.compressed_page_size)?;
995                    let buffer = self.reader.get_bytes(front.offset as u64, page_len)?;
996
997                    let (offset, header) = Self::read_page_header_len_from_bytes(
998                        &self.context,
999                        buffer.as_ref(),
1000                        *page_index,
1001                        is_dictionary_page,
1002                    )?;
1003                    let bytes = buffer.slice(offset..);
1004                    let bytes =
1005                        self.context
1006                            .decrypt_page_data(bytes, *page_index, is_dictionary_page)?;
1007
1008                    if !is_dictionary_page {
1009                        *page_index += 1;
1010                    }
1011                    decode_page(
1012                        header,
1013                        bytes,
1014                        self.physical_type,
1015                        self.decompressor.as_mut(),
1016                    )?
1017                }
1018            };
1019
1020            return Ok(Some(page));
1021        }
1022    }
1023
1024    fn peek_next_page(&mut self) -> Result<Option<PageMetadata>> {
1025        match &mut self.state {
1026            SerializedPageReaderState::Values {
1027                offset,
1028                remaining_bytes,
1029                next_page_header,
1030                page_index,
1031                require_dictionary,
1032            } => {
1033                loop {
1034                    if *remaining_bytes == 0 {
1035                        return Ok(None);
1036                    }
1037                    return if let Some(header) = next_page_header.as_ref() {
1038                        if let Ok(page_meta) = (&**header).try_into() {
1039                            Ok(Some(page_meta))
1040                        } else {
1041                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
1042                            *next_page_header = None;
1043                            continue;
1044                        }
1045                    } else {
1046                        let mut read = self.reader.get_read(*offset)?;
1047                        let (header_len, header) = Self::read_page_header_len(
1048                            &self.context,
1049                            &mut read,
1050                            *page_index,
1051                            *require_dictionary,
1052                        )?;
1053                        verify_page_header_len(header_len, *remaining_bytes)?;
1054                        *offset += header_len as u64;
1055                        *remaining_bytes -= header_len as u64;
1056                        let page_meta = if let Ok(page_meta) = (&header).try_into() {
1057                            Ok(Some(page_meta))
1058                        } else {
1059                            // For unknown page type (e.g., INDEX_PAGE), skip and read next.
1060                            continue;
1061                        };
1062                        *next_page_header = Some(Box::new(header));
1063                        page_meta
1064                    };
1065                }
1066            }
1067            SerializedPageReaderState::Pages {
1068                page_locations,
1069                dictionary_page,
1070                total_rows,
1071                page_index: _,
1072            } => {
1073                if dictionary_page.is_some() {
1074                    Ok(Some(PageMetadata {
1075                        num_rows: None,
1076                        num_levels: None,
1077                        is_dict: true,
1078                    }))
1079                } else if let Some(page) = page_locations.front() {
1080                    let next_rows = page_locations
1081                        .get(1)
1082                        .map(|x| x.first_row_index as usize)
1083                        .unwrap_or(*total_rows);
1084
1085                    Ok(Some(PageMetadata {
1086                        num_rows: Some(next_rows - page.first_row_index as usize),
1087                        num_levels: None,
1088                        is_dict: false,
1089                    }))
1090                } else {
1091                    Ok(None)
1092                }
1093            }
1094        }
1095    }
1096
1097    fn skip_next_page(&mut self) -> Result<()> {
1098        match &mut self.state {
1099            SerializedPageReaderState::Values {
1100                offset,
1101                remaining_bytes,
1102                next_page_header,
1103                page_index,
1104                require_dictionary,
1105            } => {
1106                if let Some(buffered_header) = next_page_header.take() {
1107                    verify_page_size(
1108                        buffered_header.compressed_page_size,
1109                        buffered_header.uncompressed_page_size,
1110                        *remaining_bytes,
1111                    )?;
1112                    // The next page header has already been peeked, so just advance the offset
1113                    *offset += buffered_header.compressed_page_size as u64;
1114                    *remaining_bytes -= buffered_header.compressed_page_size as u64;
1115                } else {
1116                    let mut read = self.reader.get_read(*offset)?;
1117                    let (header_len, header) = Self::read_page_header_len(
1118                        &self.context,
1119                        &mut read,
1120                        *page_index,
1121                        *require_dictionary,
1122                    )?;
1123                    verify_page_header_len(header_len, *remaining_bytes)?;
1124                    verify_page_size(
1125                        header.compressed_page_size,
1126                        header.uncompressed_page_size,
1127                        *remaining_bytes,
1128                    )?;
1129                    let data_page_size = header.compressed_page_size as u64;
1130                    *offset += header_len as u64 + data_page_size;
1131                    *remaining_bytes -= header_len as u64 + data_page_size;
1132                }
1133                if *require_dictionary {
1134                    *require_dictionary = false;
1135                } else {
1136                    *page_index += 1;
1137                }
1138                Ok(())
1139            }
1140            SerializedPageReaderState::Pages {
1141                page_locations,
1142                dictionary_page,
1143                page_index,
1144                ..
1145            } => {
1146                if dictionary_page.is_some() {
1147                    // If a dictionary page exists, consume it by taking it (sets to None)
1148                    dictionary_page.take();
1149                } else {
1150                    // If no dictionary page exists, simply pop the data page from page_locations
1151                    if page_locations.pop_front().is_some() {
1152                        *page_index += 1;
1153                    }
1154                }
1155
1156                Ok(())
1157            }
1158        }
1159    }
1160
1161    fn at_record_boundary(&mut self) -> Result<bool> {
1162        match &mut self.state {
1163            SerializedPageReaderState::Values { .. } => match self.peek_next_page()? {
1164                None => Ok(true),
1165                // V2 data pages must start at record boundaries per the parquet
1166                // spec, so the current page ends at one.
1167                Some(metadata) => Ok(metadata.num_rows.is_some()),
1168            },
1169            SerializedPageReaderState::Pages { .. } => Ok(true),
1170        }
1171    }
1172}
1173
1174#[cfg(test)]
1175mod tests {
1176    use std::collections::HashSet;
1177
1178    use bytes::Buf;
1179
1180    use crate::file::page_index::column_index::{
1181        ByteArrayColumnIndex, ColumnIndexMetaData, PrimitiveColumnIndex,
1182    };
1183    use crate::file::properties::{EnabledStatistics, WriterProperties};
1184
1185    use crate::basic::{self, BoundaryOrder, ColumnOrder, Encoding, SortOrder};
1186    use crate::column::reader::ColumnReader;
1187    use crate::data_type::private::ParquetValueType;
1188    use crate::data_type::{AsBytes, FixedLenByteArrayType, Int32Type};
1189    use crate::file::metadata::thrift::DataPageHeaderV2;
1190    use crate::file::writer::SerializedFileWriter;
1191    use crate::record::RowAccessor;
1192    use crate::schema::parser::parse_message_type;
1193    use crate::util::test_common::file_util::{get_test_file, get_test_path};
1194
1195    use super::*;
1196
1197    #[test]
1198    fn test_decode_page_invalid_offset() {
1199        let page_header = PageHeader {
1200            r#type: PageType::DATA_PAGE_V2,
1201            uncompressed_page_size: 10,
1202            compressed_page_size: 10,
1203            data_page_header: None,
1204            index_page_header: None,
1205            dictionary_page_header: None,
1206            crc: None,
1207            data_page_header_v2: Some(DataPageHeaderV2 {
1208                num_nulls: 0,
1209                num_rows: 0,
1210                num_values: 0,
1211                encoding: Encoding::PLAIN,
1212                definition_levels_byte_length: 11,
1213                repetition_levels_byte_length: 0,
1214                is_compressed: None,
1215                statistics: None,
1216            }),
1217        };
1218
1219        let buffer = Bytes::new();
1220        let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1221        assert!(
1222            err.to_string()
1223                .contains("DataPage v2 header contains implausible values")
1224        );
1225    }
1226
1227    #[test]
1228    fn test_decode_unsupported_page() {
1229        let mut page_header = PageHeader {
1230            r#type: PageType::INDEX_PAGE,
1231            uncompressed_page_size: 10,
1232            compressed_page_size: 10,
1233            data_page_header: None,
1234            index_page_header: None,
1235            dictionary_page_header: None,
1236            crc: None,
1237            data_page_header_v2: None,
1238        };
1239        let buffer = Bytes::new();
1240        let err = decode_page(page_header.clone(), buffer.clone(), Type::INT32, None).unwrap_err();
1241        assert_eq!(
1242            err.to_string(),
1243            "Parquet error: Page type INDEX_PAGE is not supported"
1244        );
1245
1246        page_header.data_page_header_v2 = Some(DataPageHeaderV2 {
1247            num_nulls: 0,
1248            num_rows: 0,
1249            num_values: 0,
1250            encoding: Encoding::PLAIN,
1251            definition_levels_byte_length: 11,
1252            repetition_levels_byte_length: 0,
1253            is_compressed: None,
1254            statistics: None,
1255        });
1256        let err = decode_page(page_header, buffer, Type::INT32, None).unwrap_err();
1257        assert!(
1258            err.to_string()
1259                .contains("DataPage v2 header contains implausible values")
1260        );
1261    }
1262
1263    #[test]
1264    fn test_cursor_and_file_has_the_same_behaviour() {
1265        let mut buf: Vec<u8> = Vec::new();
1266        get_test_file("alltypes_plain.parquet")
1267            .read_to_end(&mut buf)
1268            .unwrap();
1269        let cursor = Bytes::from(buf);
1270        let read_from_cursor = SerializedFileReader::new(cursor).unwrap();
1271
1272        let test_file = get_test_file("alltypes_plain.parquet");
1273        let read_from_file = SerializedFileReader::new(test_file).unwrap();
1274
1275        let file_iter = read_from_file.get_row_iter(None).unwrap();
1276        let cursor_iter = read_from_cursor.get_row_iter(None).unwrap();
1277
1278        for (a, b) in file_iter.zip(cursor_iter) {
1279            assert_eq!(a.unwrap(), b.unwrap())
1280        }
1281    }
1282
1283    #[test]
1284    fn test_file_reader_try_from() {
1285        // Valid file path
1286        let test_file = get_test_file("alltypes_plain.parquet");
1287        let test_path_buf = get_test_path("alltypes_plain.parquet");
1288        let test_path = test_path_buf.as_path();
1289        let test_path_str = test_path.to_str().unwrap();
1290
1291        let reader = SerializedFileReader::try_from(test_file);
1292        assert!(reader.is_ok());
1293
1294        let reader = SerializedFileReader::try_from(test_path);
1295        assert!(reader.is_ok());
1296
1297        let reader = SerializedFileReader::try_from(test_path_str);
1298        assert!(reader.is_ok());
1299
1300        let reader = SerializedFileReader::try_from(test_path_str.to_string());
1301        assert!(reader.is_ok());
1302
1303        // Invalid file path
1304        let test_path = Path::new("invalid.parquet");
1305        let test_path_str = test_path.to_str().unwrap();
1306
1307        let reader = SerializedFileReader::try_from(test_path);
1308        assert!(reader.is_err());
1309
1310        let reader = SerializedFileReader::try_from(test_path_str);
1311        assert!(reader.is_err());
1312
1313        let reader = SerializedFileReader::try_from(test_path_str.to_string());
1314        assert!(reader.is_err());
1315    }
1316
1317    #[test]
1318    fn test_file_reader_into_iter() {
1319        let path = get_test_path("alltypes_plain.parquet");
1320        let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1321        let iter = reader.into_iter();
1322        let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1323
1324        assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1325    }
1326
1327    #[test]
1328    fn test_file_reader_into_iter_project() {
1329        let path = get_test_path("alltypes_plain.parquet");
1330        let reader = SerializedFileReader::try_from(path.as_path()).unwrap();
1331        let schema = "message schema { OPTIONAL INT32 id; }";
1332        let proj = parse_message_type(schema).ok();
1333        let iter = reader.into_iter().project(proj).unwrap();
1334        let values: Vec<_> = iter.flat_map(|x| x.unwrap().get_int(0)).collect();
1335
1336        assert_eq!(values, &[4, 5, 6, 7, 2, 3, 0, 1]);
1337    }
1338
1339    #[test]
1340    fn test_reuse_file_chunk() {
1341        // This test covers the case of maintaining the correct start position in a file
1342        // stream for each column reader after initializing and moving to the next one
1343        // (without necessarily reading the entire column).
1344        let test_file = get_test_file("alltypes_plain.parquet");
1345        let reader = SerializedFileReader::new(test_file).unwrap();
1346        let row_group = reader.get_row_group(0).unwrap();
1347
1348        let mut page_readers = Vec::new();
1349        for i in 0..row_group.num_columns() {
1350            page_readers.push(row_group.get_column_page_reader(i).unwrap());
1351        }
1352
1353        // Now buffer each col reader, we do not expect any failures like:
1354        // General("underlying Thrift error: end of file")
1355        for mut page_reader in page_readers {
1356            assert!(page_reader.get_next_page().is_ok());
1357        }
1358    }
1359
1360    #[test]
1361    fn test_file_reader() {
1362        let test_file = get_test_file("alltypes_plain.parquet");
1363        let reader_result = SerializedFileReader::new(test_file);
1364        assert!(reader_result.is_ok());
1365        let reader = reader_result.unwrap();
1366
1367        // Test contents in Parquet metadata
1368        let metadata = reader.metadata();
1369        assert_eq!(metadata.num_row_groups(), 1);
1370
1371        // Test contents in file metadata
1372        let file_metadata = metadata.file_metadata();
1373        assert!(file_metadata.created_by().is_some());
1374        assert_eq!(
1375            file_metadata.created_by().unwrap(),
1376            "impala version 1.3.0-INTERNAL (build 8a48ddb1eff84592b3fc06bc6f51ec120e1fffc9)"
1377        );
1378        assert!(file_metadata.key_value_metadata().is_none());
1379        assert_eq!(file_metadata.num_rows(), 8);
1380        assert_eq!(file_metadata.version(), 1);
1381        assert_eq!(file_metadata.column_orders(), None);
1382
1383        // Test contents in row group metadata
1384        let row_group_metadata = metadata.row_group(0);
1385        assert_eq!(row_group_metadata.num_columns(), 11);
1386        assert_eq!(row_group_metadata.num_rows(), 8);
1387        assert_eq!(row_group_metadata.total_byte_size(), 671);
1388        // Check each column order
1389        for i in 0..row_group_metadata.num_columns() {
1390            assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1391        }
1392
1393        // Test row group reader
1394        let row_group_reader_result = reader.get_row_group(0);
1395        assert!(row_group_reader_result.is_ok());
1396        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1397        assert_eq!(
1398            row_group_reader.num_columns(),
1399            row_group_metadata.num_columns()
1400        );
1401        assert_eq!(
1402            row_group_reader.metadata().total_byte_size(),
1403            row_group_metadata.total_byte_size()
1404        );
1405
1406        // Test page readers
1407        // TODO: test for every column
1408        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1409        assert!(page_reader_0_result.is_ok());
1410        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1411        let mut page_count = 0;
1412        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1413            let is_expected_page = match page {
1414                Page::DictionaryPage {
1415                    buf,
1416                    num_values,
1417                    encoding,
1418                    is_sorted,
1419                } => {
1420                    assert_eq!(buf.len(), 32);
1421                    assert_eq!(num_values, 8);
1422                    assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1423                    assert!(!is_sorted);
1424                    true
1425                }
1426                Page::DataPage {
1427                    buf,
1428                    num_values,
1429                    encoding,
1430                    def_level_encoding,
1431                    rep_level_encoding,
1432                    statistics,
1433                } => {
1434                    assert_eq!(buf.len(), 11);
1435                    assert_eq!(num_values, 8);
1436                    assert_eq!(encoding, Encoding::PLAIN_DICTIONARY);
1437                    assert_eq!(def_level_encoding, Encoding::RLE);
1438                    #[allow(deprecated)]
1439                    let expected_rep_level_encoding = Encoding::BIT_PACKED;
1440                    assert_eq!(rep_level_encoding, expected_rep_level_encoding);
1441                    assert!(statistics.is_none());
1442                    true
1443                }
1444                _ => false,
1445            };
1446            assert!(is_expected_page);
1447            page_count += 1;
1448        }
1449        assert_eq!(page_count, 2);
1450    }
1451
1452    #[test]
1453    fn test_file_reader_datapage_v2() {
1454        let test_file = get_test_file("datapage_v2.snappy.parquet");
1455        let reader_result = SerializedFileReader::new(test_file);
1456        assert!(reader_result.is_ok());
1457        let reader = reader_result.unwrap();
1458
1459        // Test contents in Parquet metadata
1460        let metadata = reader.metadata();
1461        assert_eq!(metadata.num_row_groups(), 1);
1462
1463        // Test contents in file metadata
1464        let file_metadata = metadata.file_metadata();
1465        assert!(file_metadata.created_by().is_some());
1466        assert_eq!(
1467            file_metadata.created_by().unwrap(),
1468            "parquet-mr version 1.8.1 (build 4aba4dae7bb0d4edbcf7923ae1339f28fd3f7fcf)"
1469        );
1470        assert!(file_metadata.key_value_metadata().is_some());
1471        assert_eq!(
1472            file_metadata.key_value_metadata().to_owned().unwrap().len(),
1473            1
1474        );
1475
1476        assert_eq!(file_metadata.num_rows(), 5);
1477        assert_eq!(file_metadata.version(), 1);
1478        assert_eq!(file_metadata.column_orders(), None);
1479
1480        let row_group_metadata = metadata.row_group(0);
1481
1482        // Check each column order
1483        for i in 0..row_group_metadata.num_columns() {
1484            assert_eq!(file_metadata.column_order(i), ColumnOrder::UNDEFINED);
1485        }
1486
1487        // Test row group reader
1488        let row_group_reader_result = reader.get_row_group(0);
1489        assert!(row_group_reader_result.is_ok());
1490        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1491        assert_eq!(
1492            row_group_reader.num_columns(),
1493            row_group_metadata.num_columns()
1494        );
1495        assert_eq!(
1496            row_group_reader.metadata().total_byte_size(),
1497            row_group_metadata.total_byte_size()
1498        );
1499
1500        // Test page readers
1501        // TODO: test for every column
1502        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1503        assert!(page_reader_0_result.is_ok());
1504        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1505        let mut page_count = 0;
1506        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1507            let is_expected_page = match page {
1508                Page::DictionaryPage {
1509                    buf,
1510                    num_values,
1511                    encoding,
1512                    is_sorted,
1513                } => {
1514                    assert_eq!(buf.len(), 7);
1515                    assert_eq!(num_values, 1);
1516                    assert_eq!(encoding, Encoding::PLAIN);
1517                    assert!(!is_sorted);
1518                    true
1519                }
1520                Page::DataPageV2 {
1521                    buf,
1522                    num_values,
1523                    encoding,
1524                    num_nulls,
1525                    num_rows,
1526                    def_levels_byte_len,
1527                    rep_levels_byte_len,
1528                    is_compressed,
1529                    statistics,
1530                } => {
1531                    assert_eq!(buf.len(), 4);
1532                    assert_eq!(num_values, 5);
1533                    assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1534                    assert_eq!(num_nulls, 1);
1535                    assert_eq!(num_rows, 5);
1536                    assert_eq!(def_levels_byte_len, 2);
1537                    assert_eq!(rep_levels_byte_len, 0);
1538                    assert!(is_compressed);
1539                    assert!(statistics.is_none()); // page stats are no longer read
1540                    true
1541                }
1542                _ => false,
1543            };
1544            assert!(is_expected_page);
1545            page_count += 1;
1546        }
1547        assert_eq!(page_count, 2);
1548    }
1549
1550    #[test]
1551    fn test_file_reader_empty_compressed_datapage_v2() {
1552        // this file has a compressed datapage that un-compresses to 0 bytes
1553        let test_file = get_test_file("page_v2_empty_compressed.parquet");
1554        let reader_result = SerializedFileReader::new(test_file);
1555        assert!(reader_result.is_ok());
1556        let reader = reader_result.unwrap();
1557
1558        // Test contents in Parquet metadata
1559        let metadata = reader.metadata();
1560        assert_eq!(metadata.num_row_groups(), 1);
1561
1562        // Test contents in file metadata
1563        let file_metadata = metadata.file_metadata();
1564        assert!(file_metadata.created_by().is_some());
1565        assert_eq!(
1566            file_metadata.created_by().unwrap(),
1567            "parquet-cpp-arrow version 14.0.2"
1568        );
1569        assert!(file_metadata.key_value_metadata().is_some());
1570        assert_eq!(
1571            file_metadata.key_value_metadata().to_owned().unwrap().len(),
1572            1
1573        );
1574
1575        assert_eq!(file_metadata.num_rows(), 10);
1576        assert_eq!(file_metadata.version(), 2);
1577        let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1578        assert_eq!(
1579            file_metadata.column_orders(),
1580            Some(vec![expected_order].as_ref())
1581        );
1582
1583        let row_group_metadata = metadata.row_group(0);
1584
1585        // Check each column order
1586        for i in 0..row_group_metadata.num_columns() {
1587            assert_eq!(file_metadata.column_order(i), expected_order);
1588        }
1589
1590        // Test row group reader
1591        let row_group_reader_result = reader.get_row_group(0);
1592        assert!(row_group_reader_result.is_ok());
1593        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1594        assert_eq!(
1595            row_group_reader.num_columns(),
1596            row_group_metadata.num_columns()
1597        );
1598        assert_eq!(
1599            row_group_reader.metadata().total_byte_size(),
1600            row_group_metadata.total_byte_size()
1601        );
1602
1603        // Test page readers
1604        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1605        assert!(page_reader_0_result.is_ok());
1606        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1607        let mut page_count = 0;
1608        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1609            let is_expected_page = match page {
1610                Page::DictionaryPage {
1611                    buf,
1612                    num_values,
1613                    encoding,
1614                    is_sorted,
1615                } => {
1616                    assert_eq!(buf.len(), 0);
1617                    assert_eq!(num_values, 0);
1618                    assert_eq!(encoding, Encoding::PLAIN);
1619                    assert!(!is_sorted);
1620                    true
1621                }
1622                Page::DataPageV2 {
1623                    buf,
1624                    num_values,
1625                    encoding,
1626                    num_nulls,
1627                    num_rows,
1628                    def_levels_byte_len,
1629                    rep_levels_byte_len,
1630                    is_compressed,
1631                    statistics,
1632                } => {
1633                    assert_eq!(buf.len(), 3);
1634                    assert_eq!(num_values, 10);
1635                    assert_eq!(encoding, Encoding::RLE_DICTIONARY);
1636                    assert_eq!(num_nulls, 10);
1637                    assert_eq!(num_rows, 10);
1638                    assert_eq!(def_levels_byte_len, 2);
1639                    assert_eq!(rep_levels_byte_len, 0);
1640                    assert!(is_compressed);
1641                    assert!(statistics.is_none()); // page stats are no longer read
1642                    true
1643                }
1644                _ => false,
1645            };
1646            assert!(is_expected_page);
1647            page_count += 1;
1648        }
1649        assert_eq!(page_count, 2);
1650    }
1651
1652    #[test]
1653    fn test_file_reader_empty_datapage_v2() {
1654        // this file has 0 bytes compressed datapage that un-compresses to 0 bytes
1655        let test_file = get_test_file("datapage_v2_empty_datapage.snappy.parquet");
1656        let reader_result = SerializedFileReader::new(test_file);
1657        assert!(reader_result.is_ok());
1658        let reader = reader_result.unwrap();
1659
1660        // Test contents in Parquet metadata
1661        let metadata = reader.metadata();
1662        assert_eq!(metadata.num_row_groups(), 1);
1663
1664        // Test contents in file metadata
1665        let file_metadata = metadata.file_metadata();
1666        assert!(file_metadata.created_by().is_some());
1667        assert_eq!(
1668            file_metadata.created_by().unwrap(),
1669            "parquet-mr version 1.13.1 (build db4183109d5b734ec5930d870cdae161e408ddba)"
1670        );
1671        assert!(file_metadata.key_value_metadata().is_some());
1672        assert_eq!(
1673            file_metadata.key_value_metadata().to_owned().unwrap().len(),
1674            2
1675        );
1676
1677        assert_eq!(file_metadata.num_rows(), 1);
1678        assert_eq!(file_metadata.version(), 1);
1679        let expected_order = ColumnOrder::TYPE_DEFINED_ORDER(SortOrder::SIGNED);
1680        assert_eq!(
1681            file_metadata.column_orders(),
1682            Some(vec![expected_order].as_ref())
1683        );
1684
1685        let row_group_metadata = metadata.row_group(0);
1686
1687        // Check each column order
1688        for i in 0..row_group_metadata.num_columns() {
1689            assert_eq!(file_metadata.column_order(i), expected_order);
1690        }
1691
1692        // Test row group reader
1693        let row_group_reader_result = reader.get_row_group(0);
1694        assert!(row_group_reader_result.is_ok());
1695        let row_group_reader: Box<dyn RowGroupReader> = row_group_reader_result.unwrap();
1696        assert_eq!(
1697            row_group_reader.num_columns(),
1698            row_group_metadata.num_columns()
1699        );
1700        assert_eq!(
1701            row_group_reader.metadata().total_byte_size(),
1702            row_group_metadata.total_byte_size()
1703        );
1704
1705        // Test page readers
1706        let page_reader_0_result = row_group_reader.get_column_page_reader(0);
1707        assert!(page_reader_0_result.is_ok());
1708        let mut page_reader_0: Box<dyn PageReader> = page_reader_0_result.unwrap();
1709        let mut page_count = 0;
1710        while let Some(page) = page_reader_0.get_next_page().unwrap() {
1711            let is_expected_page = match page {
1712                Page::DataPageV2 {
1713                    buf,
1714                    num_values,
1715                    encoding,
1716                    num_nulls,
1717                    num_rows,
1718                    def_levels_byte_len,
1719                    rep_levels_byte_len,
1720                    is_compressed,
1721                    statistics,
1722                } => {
1723                    assert_eq!(buf.len(), 2);
1724                    assert_eq!(num_values, 1);
1725                    assert_eq!(encoding, Encoding::PLAIN);
1726                    assert_eq!(num_nulls, 1);
1727                    assert_eq!(num_rows, 1);
1728                    assert_eq!(def_levels_byte_len, 2);
1729                    assert_eq!(rep_levels_byte_len, 0);
1730                    assert!(is_compressed);
1731                    assert!(statistics.is_none());
1732                    true
1733                }
1734                _ => false,
1735            };
1736            assert!(is_expected_page);
1737            page_count += 1;
1738        }
1739        assert_eq!(page_count, 1);
1740    }
1741
1742    fn get_serialized_page_reader<R: ChunkReader>(
1743        file_reader: &SerializedFileReader<R>,
1744        row_group: usize,
1745        column: usize,
1746    ) -> Result<SerializedPageReader<R>> {
1747        let row_group = {
1748            let row_group_metadata = file_reader.metadata.row_group(row_group);
1749            let props = Arc::clone(&file_reader.props);
1750            let f = Arc::clone(&file_reader.chunk_reader);
1751            SerializedRowGroupReader::new(
1752                f,
1753                row_group_metadata,
1754                file_reader
1755                    .metadata
1756                    .offset_index()
1757                    .map(|x| x[row_group].as_slice()),
1758                props,
1759            )?
1760        };
1761
1762        let col = row_group.metadata.column(column);
1763
1764        let page_locations = row_group
1765            .offset_index
1766            .map(|x| x[column].page_locations.clone());
1767
1768        let props = Arc::clone(&row_group.props);
1769        SerializedPageReader::new_with_properties(
1770            Arc::clone(&row_group.chunk_reader),
1771            col,
1772            usize::try_from(row_group.metadata.num_rows())?,
1773            page_locations,
1774            props,
1775        )
1776    }
1777
1778    #[test]
1779    fn test_peek_next_page_offset_matches_actual() -> Result<()> {
1780        let test_file = get_test_file("alltypes_plain.parquet");
1781        let reader = SerializedFileReader::new(test_file)?;
1782
1783        let mut offset_set = HashSet::new();
1784        let num_row_groups = reader.metadata.num_row_groups();
1785        for row_group in 0..num_row_groups {
1786            let num_columns = reader.metadata.row_group(row_group).num_columns();
1787            for column in 0..num_columns {
1788                let mut page_reader = get_serialized_page_reader(&reader, row_group, column)?;
1789
1790                while let Ok(Some(page_offset)) = page_reader.peek_next_page_offset() {
1791                    match &page_reader.state {
1792                        SerializedPageReaderState::Pages {
1793                            page_locations,
1794                            dictionary_page,
1795                            ..
1796                        } => {
1797                            if let Some(page) = dictionary_page {
1798                                assert_eq!(page.offset as u64, page_offset);
1799                            } else if let Some(page) = page_locations.front() {
1800                                assert_eq!(page.offset as u64, page_offset);
1801                            } else {
1802                                unreachable!()
1803                            }
1804                        }
1805                        SerializedPageReaderState::Values {
1806                            offset,
1807                            next_page_header,
1808                            ..
1809                        } => {
1810                            assert!(next_page_header.is_some());
1811                            assert_eq!(*offset, page_offset);
1812                        }
1813                    }
1814                    let page = page_reader.get_next_page()?;
1815                    assert!(page.is_some());
1816                    let newly_inserted = offset_set.insert(page_offset);
1817                    assert!(newly_inserted);
1818                }
1819            }
1820        }
1821
1822        Ok(())
1823    }
1824
1825    #[test]
1826    fn test_page_iterator() {
1827        let file = get_test_file("alltypes_plain.parquet");
1828        let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1829
1830        let mut page_iterator = FilePageIterator::new(0, file_reader.clone()).unwrap();
1831
1832        // read first page
1833        let page = page_iterator.next();
1834        assert!(page.is_some());
1835        assert!(page.unwrap().is_ok());
1836
1837        // reach end of file
1838        let page = page_iterator.next();
1839        assert!(page.is_none());
1840
1841        let row_group_indices = Box::new(0..1);
1842        let mut page_iterator =
1843            FilePageIterator::with_row_groups(0, row_group_indices, file_reader).unwrap();
1844
1845        // read first page
1846        let page = page_iterator.next();
1847        assert!(page.is_some());
1848        assert!(page.unwrap().is_ok());
1849
1850        // reach end of file
1851        let page = page_iterator.next();
1852        assert!(page.is_none());
1853    }
1854
1855    #[test]
1856    fn test_file_reader_key_value_metadata() {
1857        let file = get_test_file("binary.parquet");
1858        let file_reader = Arc::new(SerializedFileReader::new(file).unwrap());
1859
1860        let metadata = file_reader
1861            .metadata
1862            .file_metadata()
1863            .key_value_metadata()
1864            .unwrap();
1865
1866        assert_eq!(metadata.len(), 3);
1867
1868        assert_eq!(metadata[0].key, "parquet.proto.descriptor");
1869
1870        assert_eq!(metadata[1].key, "writer.model.name");
1871        assert_eq!(metadata[1].value, Some("protobuf".to_owned()));
1872
1873        assert_eq!(metadata[2].key, "parquet.proto.class");
1874        assert_eq!(metadata[2].value, Some("foo.baz.Foobaz$Event".to_owned()));
1875    }
1876
1877    #[test]
1878    fn test_file_reader_optional_metadata() {
1879        // file with optional metadata: bloom filters, encoding stats, column index and offset index.
1880        let file = get_test_file("data_index_bloom_encoding_stats.parquet");
1881        let options = ReadOptionsBuilder::new()
1882            .with_encoding_stats_as_mask(false)
1883            .build();
1884        let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1885
1886        let row_group_metadata = file_reader.metadata.row_group(0);
1887        let col0_metadata = row_group_metadata.column(0);
1888
1889        // test optional bloom filter offset
1890        assert_eq!(col0_metadata.bloom_filter_offset().unwrap(), 192);
1891
1892        // test page encoding stats
1893        let page_encoding_stats = &col0_metadata.page_encoding_stats().unwrap()[0];
1894
1895        assert_eq!(page_encoding_stats.page_type, basic::PageType::DATA_PAGE);
1896        assert_eq!(page_encoding_stats.encoding, Encoding::PLAIN);
1897        assert_eq!(page_encoding_stats.count, 1);
1898
1899        // test optional column index offset
1900        assert_eq!(col0_metadata.column_index_offset().unwrap(), 156);
1901        assert_eq!(col0_metadata.column_index_length().unwrap(), 25);
1902
1903        // test optional offset index offset
1904        assert_eq!(col0_metadata.offset_index_offset().unwrap(), 181);
1905        assert_eq!(col0_metadata.offset_index_length().unwrap(), 11);
1906    }
1907
1908    #[test]
1909    fn test_file_reader_page_stats_mask() {
1910        let file = get_test_file("alltypes_tiny_pages.parquet");
1911        let options = ReadOptionsBuilder::new()
1912            .with_encoding_stats_as_mask(true)
1913            .build();
1914        let file_reader = Arc::new(SerializedFileReader::new_with_options(file, options).unwrap());
1915
1916        let row_group_metadata = file_reader.metadata.row_group(0);
1917
1918        // test page encoding stats
1919        let page_encoding_stats = row_group_metadata
1920            .column(0)
1921            .page_encoding_stats_mask()
1922            .unwrap();
1923        assert!(page_encoding_stats.is_only(Encoding::PLAIN));
1924        let page_encoding_stats = row_group_metadata
1925            .column(2)
1926            .page_encoding_stats_mask()
1927            .unwrap();
1928        assert!(page_encoding_stats.is_only(Encoding::PLAIN_DICTIONARY));
1929    }
1930
1931    #[test]
1932    fn test_file_reader_page_stats_skipped() {
1933        let file = get_test_file("alltypes_tiny_pages.parquet");
1934
1935        // test skipping all
1936        let options = ReadOptionsBuilder::new()
1937            .with_encoding_stats_policy(ParquetStatisticsPolicy::SkipAll)
1938            .with_column_stats_policy(ParquetStatisticsPolicy::SkipAll)
1939            .build();
1940        let file_reader = Arc::new(
1941            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1942        );
1943
1944        let row_group_metadata = file_reader.metadata.row_group(0);
1945        for column in row_group_metadata.columns() {
1946            assert!(column.page_encoding_stats().is_none());
1947            assert!(column.page_encoding_stats_mask().is_none());
1948            assert!(column.statistics().is_none());
1949        }
1950
1951        // test skipping all but one column
1952        let options = ReadOptionsBuilder::new()
1953            .with_encoding_stats_as_mask(true)
1954            .with_encoding_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1955            .with_column_stats_policy(ParquetStatisticsPolicy::skip_except(&[0]))
1956            .build();
1957        let file_reader = Arc::new(
1958            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1959        );
1960
1961        let row_group_metadata = file_reader.metadata.row_group(0);
1962        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
1963            assert!(column.page_encoding_stats().is_none());
1964            assert_eq!(column.page_encoding_stats_mask().is_some(), idx == 0);
1965            assert_eq!(column.statistics().is_some(), idx == 0);
1966        }
1967    }
1968
1969    #[test]
1970    fn test_file_reader_size_stats_skipped() {
1971        let file = get_test_file("repeated_primitive_no_list.parquet");
1972
1973        // test skipping all
1974        let options = ReadOptionsBuilder::new()
1975            .with_size_stats_policy(ParquetStatisticsPolicy::SkipAll)
1976            .build();
1977        let file_reader = Arc::new(
1978            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1979        );
1980
1981        let row_group_metadata = file_reader.metadata.row_group(0);
1982        for column in row_group_metadata.columns() {
1983            assert!(column.repetition_level_histogram().is_none());
1984            assert!(column.definition_level_histogram().is_none());
1985            assert!(column.unencoded_byte_array_data_bytes().is_none());
1986        }
1987
1988        // test skipping all but one column
1989        let options = ReadOptionsBuilder::new()
1990            .with_encoding_stats_as_mask(true)
1991            .with_size_stats_policy(ParquetStatisticsPolicy::skip_except(&[1]))
1992            .build();
1993        let file_reader = Arc::new(
1994            SerializedFileReader::new_with_options(file.try_clone().unwrap(), options).unwrap(),
1995        );
1996
1997        let row_group_metadata = file_reader.metadata.row_group(0);
1998        for (idx, column) in row_group_metadata.columns().iter().enumerate() {
1999            assert_eq!(column.repetition_level_histogram().is_some(), idx == 1);
2000            assert_eq!(column.definition_level_histogram().is_some(), idx == 1);
2001            assert_eq!(column.unencoded_byte_array_data_bytes().is_some(), idx == 1);
2002        }
2003    }
2004
2005    #[test]
2006    fn test_file_reader_with_no_filter() -> Result<()> {
2007        let test_file = get_test_file("alltypes_plain.parquet");
2008        let origin_reader = SerializedFileReader::new(test_file)?;
2009        // test initial number of row groups
2010        let metadata = origin_reader.metadata();
2011        assert_eq!(metadata.num_row_groups(), 1);
2012        Ok(())
2013    }
2014
2015    #[test]
2016    fn test_file_reader_filter_row_groups_with_predicate() -> Result<()> {
2017        let test_file = get_test_file("alltypes_plain.parquet");
2018        let read_options = ReadOptionsBuilder::new()
2019            .with_predicate(Box::new(|_, _| false))
2020            .build();
2021        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2022        let metadata = reader.metadata();
2023        assert_eq!(metadata.num_row_groups(), 0);
2024        Ok(())
2025    }
2026
2027    #[test]
2028    fn test_file_reader_filter_row_groups_with_range() -> Result<()> {
2029        let test_file = get_test_file("alltypes_plain.parquet");
2030        let origin_reader = SerializedFileReader::new(test_file)?;
2031        // test initial number of row groups
2032        let metadata = origin_reader.metadata();
2033        assert_eq!(metadata.num_row_groups(), 1);
2034        let mid = get_midpoint_offset(metadata.row_group(0));
2035
2036        let test_file = get_test_file("alltypes_plain.parquet");
2037        let read_options = ReadOptionsBuilder::new().with_range(0, mid + 1).build();
2038        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2039        let metadata = reader.metadata();
2040        assert_eq!(metadata.num_row_groups(), 1);
2041
2042        let test_file = get_test_file("alltypes_plain.parquet");
2043        let read_options = ReadOptionsBuilder::new().with_range(0, mid).build();
2044        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2045        let metadata = reader.metadata();
2046        assert_eq!(metadata.num_row_groups(), 0);
2047        Ok(())
2048    }
2049
2050    #[test]
2051    fn test_file_reader_filter_row_groups_and_range() -> Result<()> {
2052        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2053        let origin_reader = SerializedFileReader::new(test_file)?;
2054        let metadata = origin_reader.metadata();
2055        let mid = get_midpoint_offset(metadata.row_group(0));
2056
2057        // true, true predicate
2058        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2059        let read_options = ReadOptionsBuilder::new()
2060            .with_page_index()
2061            .with_predicate(Box::new(|_, _| true))
2062            .with_range(mid, mid + 1)
2063            .build();
2064        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2065        let metadata = reader.metadata();
2066        assert_eq!(metadata.num_row_groups(), 1);
2067        assert_eq!(metadata.column_index().unwrap().len(), 1);
2068        assert_eq!(metadata.offset_index().unwrap().len(), 1);
2069
2070        // true, false predicate
2071        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2072        let read_options = ReadOptionsBuilder::new()
2073            .with_page_index()
2074            .with_predicate(Box::new(|_, _| true))
2075            .with_range(0, mid)
2076            .build();
2077        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2078        let metadata = reader.metadata();
2079        assert_eq!(metadata.num_row_groups(), 0);
2080        assert!(metadata.column_index().is_none());
2081        assert!(metadata.offset_index().is_none());
2082
2083        // false, true predicate
2084        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2085        let read_options = ReadOptionsBuilder::new()
2086            .with_page_index()
2087            .with_predicate(Box::new(|_, _| false))
2088            .with_range(mid, mid + 1)
2089            .build();
2090        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2091        let metadata = reader.metadata();
2092        assert_eq!(metadata.num_row_groups(), 0);
2093        assert!(metadata.column_index().is_none());
2094        assert!(metadata.offset_index().is_none());
2095
2096        // false, false predicate
2097        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2098        let read_options = ReadOptionsBuilder::new()
2099            .with_page_index()
2100            .with_predicate(Box::new(|_, _| false))
2101            .with_range(0, mid)
2102            .build();
2103        let reader = SerializedFileReader::new_with_options(test_file, read_options)?;
2104        let metadata = reader.metadata();
2105        assert_eq!(metadata.num_row_groups(), 0);
2106        assert!(metadata.column_index().is_none());
2107        assert!(metadata.offset_index().is_none());
2108        Ok(())
2109    }
2110
2111    #[test]
2112    fn test_file_reader_invalid_metadata() {
2113        let data = [
2114            255, 172, 1, 0, 50, 82, 65, 73, 1, 0, 0, 0, 169, 168, 168, 162, 87, 255, 16, 0, 0, 0,
2115            80, 65, 82, 49,
2116        ];
2117        let ret = SerializedFileReader::new(Bytes::copy_from_slice(&data));
2118        assert_eq!(
2119            ret.err().unwrap().to_string(),
2120            "Parquet error: Expected list element type of Struct but got List"
2121        );
2122    }
2123
2124    #[test]
2125    // Use java parquet-tools get below pageIndex info
2126    // !```
2127    // parquet-tools column-index ./data_index_bloom_encoding_stats.parquet
2128    // row group 0:
2129    // column index for column String:
2130    // Boundary order: ASCENDING
2131    // page-0  :
2132    // null count                 min                                  max
2133    // 0                          Hello                                today
2134    //
2135    // offset index for column String:
2136    // page-0   :
2137    // offset   compressed size       first row index
2138    // 4               152                     0
2139    ///```
2140    //
2141    fn test_page_index_reader() {
2142        let test_file = get_test_file("data_index_bloom_encoding_stats.parquet");
2143        let builder = ReadOptionsBuilder::new();
2144        //enable read page index
2145        let options = builder.with_page_index().build();
2146        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2147        let reader = reader_result.unwrap();
2148
2149        // Test contents in Parquet metadata
2150        let metadata = reader.metadata();
2151        assert_eq!(metadata.num_row_groups(), 1);
2152
2153        let column_index = metadata.column_index().unwrap();
2154
2155        // only one row group
2156        assert_eq!(column_index.len(), 1);
2157        let index = if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][0] {
2158            index
2159        } else {
2160            unreachable!()
2161        };
2162
2163        assert_eq!(index.boundary_order, BoundaryOrder::ASCENDING);
2164
2165        //only one page group
2166        assert_eq!(index.num_pages(), 1);
2167
2168        let min = index.min_value(0).unwrap();
2169        let max = index.max_value(0).unwrap();
2170        assert_eq!(b"Hello", min.as_bytes());
2171        assert_eq!(b"today", max.as_bytes());
2172
2173        let offset_indexes = metadata.offset_index().unwrap();
2174        // only one row group
2175        assert_eq!(offset_indexes.len(), 1);
2176        let offset_index = &offset_indexes[0];
2177        let page_offset = &offset_index[0].page_locations()[0];
2178
2179        assert_eq!(4, page_offset.offset);
2180        assert_eq!(152, page_offset.compressed_page_size);
2181        assert_eq!(0, page_offset.first_row_index);
2182    }
2183
2184    #[test]
2185    fn test_page_index_reader_all_type() {
2186        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2187        let builder = ReadOptionsBuilder::new();
2188        //enable read page index
2189        let options = builder.with_page_index().build();
2190        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2191        let reader = reader_result.unwrap();
2192
2193        // Test contents in Parquet metadata
2194        let metadata = reader.metadata();
2195        assert_eq!(metadata.num_row_groups(), 1);
2196
2197        let column_index = metadata.column_index().unwrap();
2198        let row_group_offset_indexes = &metadata.offset_index().unwrap()[0];
2199
2200        // only one row group
2201        assert_eq!(column_index.len(), 1);
2202        let row_group_metadata = metadata.row_group(0);
2203
2204        //col0->id: INT32 UNCOMPRESSED DO:0 FPO:4 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 7299, num_nulls: 0]
2205        assert!(!&column_index[0][0].is_sorted());
2206        let boundary_order = &column_index[0][0].get_boundary_order();
2207        assert!(boundary_order.is_some());
2208        matches!(boundary_order.unwrap(), BoundaryOrder::UNORDERED);
2209        if let ColumnIndexMetaData::INT32(index) = &column_index[0][0] {
2210            check_native_page_index(
2211                index,
2212                325,
2213                get_row_group_min_max_bytes(row_group_metadata, 0),
2214                BoundaryOrder::UNORDERED,
2215            );
2216            assert_eq!(row_group_offset_indexes[0].page_locations.len(), 325);
2217        } else {
2218            unreachable!()
2219        };
2220        //col1->bool_col:BOOLEAN UNCOMPRESSED DO:0 FPO:37329 SZ:3022/3022/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: false, max: true, num_nulls: 0]
2221        assert!(&column_index[0][1].is_sorted());
2222        if let ColumnIndexMetaData::BOOLEAN(index) = &column_index[0][1] {
2223            assert_eq!(index.num_pages(), 82);
2224            assert_eq!(row_group_offset_indexes[1].page_locations.len(), 82);
2225        } else {
2226            unreachable!()
2227        };
2228        //col2->tinyint_col: INT32 UNCOMPRESSED DO:0 FPO:40351 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2229        assert!(&column_index[0][2].is_sorted());
2230        if let ColumnIndexMetaData::INT32(index) = &column_index[0][2] {
2231            check_native_page_index(
2232                index,
2233                325,
2234                get_row_group_min_max_bytes(row_group_metadata, 2),
2235                BoundaryOrder::ASCENDING,
2236            );
2237            assert_eq!(row_group_offset_indexes[2].page_locations.len(), 325);
2238        } else {
2239            unreachable!()
2240        };
2241        //col4->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2242        assert!(&column_index[0][3].is_sorted());
2243        if let ColumnIndexMetaData::INT32(index) = &column_index[0][3] {
2244            check_native_page_index(
2245                index,
2246                325,
2247                get_row_group_min_max_bytes(row_group_metadata, 3),
2248                BoundaryOrder::ASCENDING,
2249            );
2250            assert_eq!(row_group_offset_indexes[3].page_locations.len(), 325);
2251        } else {
2252            unreachable!()
2253        };
2254        //col5->smallint_col: INT32 UNCOMPRESSED DO:0 FPO:77676 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2255        assert!(&column_index[0][4].is_sorted());
2256        if let ColumnIndexMetaData::INT32(index) = &column_index[0][4] {
2257            check_native_page_index(
2258                index,
2259                325,
2260                get_row_group_min_max_bytes(row_group_metadata, 4),
2261                BoundaryOrder::ASCENDING,
2262            );
2263            assert_eq!(row_group_offset_indexes[4].page_locations.len(), 325);
2264        } else {
2265            unreachable!()
2266        };
2267        //col6->bigint_col: INT64 UNCOMPRESSED DO:0 FPO:152326 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 90, num_nulls: 0]
2268        assert!(!&column_index[0][5].is_sorted());
2269        if let ColumnIndexMetaData::INT64(index) = &column_index[0][5] {
2270            check_native_page_index(
2271                index,
2272                528,
2273                get_row_group_min_max_bytes(row_group_metadata, 5),
2274                BoundaryOrder::UNORDERED,
2275            );
2276            assert_eq!(row_group_offset_indexes[5].page_locations.len(), 528);
2277        } else {
2278            unreachable!()
2279        };
2280        //col7->float_col: FLOAT UNCOMPRESSED DO:0 FPO:223924 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 9.9, num_nulls: 0]
2281        assert!(&column_index[0][6].is_sorted());
2282        if let ColumnIndexMetaData::FLOAT(index) = &column_index[0][6] {
2283            check_native_page_index(
2284                index,
2285                325,
2286                get_row_group_min_max_bytes(row_group_metadata, 6),
2287                BoundaryOrder::ASCENDING,
2288            );
2289            assert_eq!(row_group_offset_indexes[6].page_locations.len(), 325);
2290        } else {
2291            unreachable!()
2292        };
2293        //col8->double_col: DOUBLE UNCOMPRESSED DO:0 FPO:261249 SZ:71598/71598/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: -0.0, max: 90.89999999999999, num_nulls: 0]
2294        assert!(!&column_index[0][7].is_sorted());
2295        if let ColumnIndexMetaData::DOUBLE(index) = &column_index[0][7] {
2296            check_native_page_index(
2297                index,
2298                528,
2299                get_row_group_min_max_bytes(row_group_metadata, 7),
2300                BoundaryOrder::UNORDERED,
2301            );
2302            assert_eq!(row_group_offset_indexes[7].page_locations.len(), 528);
2303        } else {
2304            unreachable!()
2305        };
2306        //col9->date_string_col: BINARY UNCOMPRESSED DO:0 FPO:332847 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 01/01/09, max: 12/31/10, num_nulls: 0]
2307        assert!(!&column_index[0][8].is_sorted());
2308        if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][8] {
2309            check_byte_array_page_index(
2310                index,
2311                974,
2312                get_row_group_min_max_bytes(row_group_metadata, 8),
2313                BoundaryOrder::UNORDERED,
2314            );
2315            assert_eq!(row_group_offset_indexes[8].page_locations.len(), 974);
2316        } else {
2317            unreachable!()
2318        };
2319        //col10->string_col: BINARY UNCOMPRESSED DO:0 FPO:444795 SZ:45298/45298/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 0, max: 9, num_nulls: 0]
2320        assert!(&column_index[0][9].is_sorted());
2321        if let ColumnIndexMetaData::BYTE_ARRAY(index) = &column_index[0][9] {
2322            check_byte_array_page_index(
2323                index,
2324                352,
2325                get_row_group_min_max_bytes(row_group_metadata, 9),
2326                BoundaryOrder::ASCENDING,
2327            );
2328            assert_eq!(row_group_offset_indexes[9].page_locations.len(), 352);
2329        } else {
2330            unreachable!()
2331        };
2332        //col11->timestamp_col: INT96 UNCOMPRESSED DO:0 FPO:490093 SZ:111948/111948/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[num_nulls: 0, min/max not defined]
2333        //Notice: min_max values for each page for this col not exits.
2334        assert!(!&column_index[0][10].is_sorted());
2335        if column_index[0][10] == ColumnIndexMetaData::NONE {
2336            assert_eq!(row_group_offset_indexes[10].page_locations.len(), 974);
2337        } else {
2338            unreachable!()
2339        };
2340        //col12->year: INT32 UNCOMPRESSED DO:0 FPO:602041 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 2009, max: 2010, num_nulls: 0]
2341        assert!(&column_index[0][11].is_sorted());
2342        if let ColumnIndexMetaData::INT32(index) = &column_index[0][11] {
2343            check_native_page_index(
2344                index,
2345                325,
2346                get_row_group_min_max_bytes(row_group_metadata, 11),
2347                BoundaryOrder::ASCENDING,
2348            );
2349            assert_eq!(row_group_offset_indexes[11].page_locations.len(), 325);
2350        } else {
2351            unreachable!()
2352        };
2353        //col13->month: INT32 UNCOMPRESSED DO:0 FPO:639366 SZ:37325/37325/1.00 VC:7300 ENC:BIT_PACKED,RLE,PLAIN ST:[min: 1, max: 12, num_nulls: 0]
2354        assert!(!&column_index[0][12].is_sorted());
2355        if let ColumnIndexMetaData::INT32(index) = &column_index[0][12] {
2356            check_native_page_index(
2357                index,
2358                325,
2359                get_row_group_min_max_bytes(row_group_metadata, 12),
2360                BoundaryOrder::UNORDERED,
2361            );
2362            assert_eq!(row_group_offset_indexes[12].page_locations.len(), 325);
2363        } else {
2364            unreachable!()
2365        };
2366    }
2367
2368    fn check_native_page_index<T: ParquetValueType>(
2369        row_group_index: &PrimitiveColumnIndex<T>,
2370        page_size: usize,
2371        min_max: (&[u8], &[u8]),
2372        boundary_order: BoundaryOrder,
2373    ) {
2374        assert_eq!(row_group_index.num_pages() as usize, page_size);
2375        assert_eq!(row_group_index.boundary_order, boundary_order);
2376        assert!(row_group_index.min_values().iter().all(|x| {
2377            x >= &T::try_from_le_slice(min_max.0).unwrap()
2378                && x <= &T::try_from_le_slice(min_max.1).unwrap()
2379        }));
2380    }
2381
2382    fn check_byte_array_page_index(
2383        row_group_index: &ByteArrayColumnIndex,
2384        page_size: usize,
2385        min_max: (&[u8], &[u8]),
2386        boundary_order: BoundaryOrder,
2387    ) {
2388        assert_eq!(row_group_index.num_pages() as usize, page_size);
2389        assert_eq!(row_group_index.boundary_order, boundary_order);
2390        for i in 0..row_group_index.num_pages() as usize {
2391            let x = row_group_index.min_value(i).unwrap();
2392            assert!(x >= min_max.0 && x <= min_max.1);
2393        }
2394    }
2395
2396    fn get_row_group_min_max_bytes(r: &RowGroupMetaData, col_num: usize) -> (&[u8], &[u8]) {
2397        let statistics = r.column(col_num).statistics().unwrap();
2398        (
2399            statistics.min_bytes_opt().unwrap_or_default(),
2400            statistics.max_bytes_opt().unwrap_or_default(),
2401        )
2402    }
2403
2404    #[test]
2405    fn test_skip_next_page_with_dictionary_page() {
2406        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2407        let builder = ReadOptionsBuilder::new();
2408        // enable read page index
2409        let options = builder.with_page_index().build();
2410        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2411        let reader = reader_result.unwrap();
2412
2413        let row_group_reader = reader.get_row_group(0).unwrap();
2414
2415        // use 'string_col', Boundary order: UNORDERED, total 352 data pages and 1 dictionary page.
2416        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2417
2418        let mut vec = vec![];
2419
2420        // Step 1: Peek and ensure dictionary page is correctly identified
2421        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2422        assert!(meta.is_dict);
2423
2424        // Step 2: Call skip_next_page to skip the dictionary page
2425        column_page_reader.skip_next_page().unwrap();
2426
2427        // Step 3: Read the next data page after skipping the dictionary page
2428        let page = column_page_reader.get_next_page().unwrap().unwrap();
2429        assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2430
2431        // Step 4: Continue reading remaining data pages and verify correctness
2432        for _i in 0..351 {
2433            // 352 total pages, 1 dictionary page is skipped
2434            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2435            assert!(!meta.is_dict); // Verify no dictionary page here
2436            vec.push(meta);
2437
2438            let page = column_page_reader.get_next_page().unwrap().unwrap();
2439            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2440        }
2441
2442        // Step 5: Check if all pages are read
2443        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2444        assert!(column_page_reader.get_next_page().unwrap().is_none());
2445
2446        // Step 6: Verify the number of data pages read (should be 351 data pages)
2447        assert_eq!(vec.len(), 351);
2448    }
2449
2450    #[test]
2451    fn test_skip_page_with_offset_index() {
2452        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2453        let builder = ReadOptionsBuilder::new();
2454        //enable read page index
2455        let options = builder.with_page_index().build();
2456        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2457        let reader = reader_result.unwrap();
2458
2459        let row_group_reader = reader.get_row_group(0).unwrap();
2460
2461        //use 'int_col', Boundary order: ASCENDING, total 325 pages.
2462        let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2463
2464        let mut vec = vec![];
2465
2466        for i in 0..325 {
2467            if i % 2 == 0 {
2468                vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2469            } else {
2470                column_page_reader.skip_next_page().unwrap();
2471            }
2472        }
2473        //check read all pages.
2474        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2475        assert!(column_page_reader.get_next_page().unwrap().is_none());
2476
2477        assert_eq!(vec.len(), 163);
2478    }
2479
2480    #[test]
2481    fn test_skip_page_without_offset_index() {
2482        let test_file = get_test_file("alltypes_tiny_pages_plain.parquet");
2483
2484        // use default SerializedFileReader without read offsetIndex
2485        let reader_result = SerializedFileReader::new(test_file);
2486        let reader = reader_result.unwrap();
2487
2488        let row_group_reader = reader.get_row_group(0).unwrap();
2489
2490        //use 'int_col', Boundary order: ASCENDING, total 325 pages.
2491        let mut column_page_reader = row_group_reader.get_column_page_reader(4).unwrap();
2492
2493        let mut vec = vec![];
2494
2495        for i in 0..325 {
2496            if i % 2 == 0 {
2497                vec.push(column_page_reader.get_next_page().unwrap().unwrap());
2498            } else {
2499                column_page_reader.peek_next_page().unwrap().unwrap();
2500                column_page_reader.skip_next_page().unwrap();
2501            }
2502        }
2503        //check read all pages.
2504        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2505        assert!(column_page_reader.get_next_page().unwrap().is_none());
2506
2507        assert_eq!(vec.len(), 163);
2508    }
2509
2510    #[test]
2511    fn test_peek_page_with_dictionary_page() {
2512        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2513        let builder = ReadOptionsBuilder::new();
2514        //enable read page index
2515        let options = builder.with_page_index().build();
2516        let reader_result = SerializedFileReader::new_with_options(test_file, options);
2517        let reader = reader_result.unwrap();
2518        let row_group_reader = reader.get_row_group(0).unwrap();
2519
2520        //use 'string_col', Boundary order: UNORDERED, total 352 data ages and 1 dictionary page.
2521        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2522
2523        let mut vec = vec![];
2524
2525        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2526        assert!(meta.is_dict);
2527        let page = column_page_reader.get_next_page().unwrap().unwrap();
2528        assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2529
2530        for i in 0..352 {
2531            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2532            // have checked with `parquet-tools column-index   -c string_col  ./alltypes_tiny_pages.parquet`
2533            // page meta has two scenarios(21, 20) of num_rows expect last page has 11 rows.
2534            if i != 351 {
2535                assert!((meta.num_rows == Some(21)) || (meta.num_rows == Some(20)));
2536            } else {
2537                // last page first row index is 7290, total row count is 7300
2538                // because first row start with zero, last page row count should be 10.
2539                assert_eq!(meta.num_rows, Some(10));
2540            }
2541            assert!(!meta.is_dict);
2542            vec.push(meta);
2543            let page = column_page_reader.get_next_page().unwrap().unwrap();
2544            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2545        }
2546
2547        //check read all pages.
2548        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2549        assert!(column_page_reader.get_next_page().unwrap().is_none());
2550
2551        assert_eq!(vec.len(), 352);
2552    }
2553
2554    #[test]
2555    fn test_peek_page_with_dictionary_page_without_offset_index() {
2556        let test_file = get_test_file("alltypes_tiny_pages.parquet");
2557
2558        let reader_result = SerializedFileReader::new(test_file);
2559        let reader = reader_result.unwrap();
2560        let row_group_reader = reader.get_row_group(0).unwrap();
2561
2562        //use 'string_col', Boundary order: UNORDERED, total 352 data ages and 1 dictionary page.
2563        let mut column_page_reader = row_group_reader.get_column_page_reader(9).unwrap();
2564
2565        let mut vec = vec![];
2566
2567        let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2568        assert!(meta.is_dict);
2569        let page = column_page_reader.get_next_page().unwrap().unwrap();
2570        assert!(matches!(page.page_type(), basic::PageType::DICTIONARY_PAGE));
2571
2572        for i in 0..352 {
2573            let meta = column_page_reader.peek_next_page().unwrap().unwrap();
2574            // have checked with `parquet-tools column-index   -c string_col  ./alltypes_tiny_pages.parquet`
2575            // page meta has two scenarios(21, 20) of num_rows expect last page has 11 rows.
2576            if i != 351 {
2577                assert!((meta.num_levels == Some(21)) || (meta.num_levels == Some(20)));
2578            } else {
2579                // last page first row index is 7290, total row count is 7300
2580                // because first row start with zero, last page row count should be 10.
2581                assert_eq!(meta.num_levels, Some(10));
2582            }
2583            assert!(!meta.is_dict);
2584            vec.push(meta);
2585            let page = column_page_reader.get_next_page().unwrap().unwrap();
2586            assert!(matches!(page.page_type(), basic::PageType::DATA_PAGE));
2587        }
2588
2589        //check read all pages.
2590        assert!(column_page_reader.peek_next_page().unwrap().is_none());
2591        assert!(column_page_reader.get_next_page().unwrap().is_none());
2592
2593        assert_eq!(vec.len(), 352);
2594    }
2595
2596    #[test]
2597    fn test_fixed_length_index() {
2598        let message_type = "
2599        message test_schema {
2600          OPTIONAL FIXED_LEN_BYTE_ARRAY (11) value (DECIMAL(25,2));
2601        }
2602        ";
2603
2604        let schema = parse_message_type(message_type).unwrap();
2605        let mut out = Vec::with_capacity(1024);
2606        let mut writer =
2607            SerializedFileWriter::new(&mut out, Arc::new(schema), Default::default()).unwrap();
2608
2609        let mut r = writer.next_row_group().unwrap();
2610        let mut c = r.next_column().unwrap().unwrap();
2611        c.typed::<FixedLenByteArrayType>()
2612            .write_batch(
2613                &[vec![0; 11].into(), vec![5; 11].into(), vec![3; 11].into()],
2614                Some(&[1, 1, 0, 1]),
2615                None,
2616            )
2617            .unwrap();
2618        c.close().unwrap();
2619        r.close().unwrap();
2620        writer.close().unwrap();
2621
2622        let b = Bytes::from(out);
2623        let options = ReadOptionsBuilder::new().with_page_index().build();
2624        let reader = SerializedFileReader::new_with_options(b, options).unwrap();
2625        let index = reader.metadata().column_index().unwrap();
2626
2627        // 1 row group
2628        assert_eq!(index.len(), 1);
2629        let c = &index[0];
2630        // 1 column
2631        assert_eq!(c.len(), 1);
2632
2633        match &c[0] {
2634            ColumnIndexMetaData::FIXED_LEN_BYTE_ARRAY(v) => {
2635                assert_eq!(v.num_pages(), 1);
2636                assert_eq!(v.null_count(0).unwrap(), 1);
2637                assert_eq!(v.min_value(0).unwrap(), &[0; 11]);
2638                assert_eq!(v.max_value(0).unwrap(), &[5; 11]);
2639            }
2640            _ => unreachable!(),
2641        }
2642    }
2643
2644    #[test]
2645    fn test_multi_gz() {
2646        let file = get_test_file("concatenated_gzip_members.parquet");
2647        let reader = SerializedFileReader::new(file).unwrap();
2648        let row_group_reader = reader.get_row_group(0).unwrap();
2649        match row_group_reader.get_column_reader(0).unwrap() {
2650            ColumnReader::Int64ColumnReader(mut reader) => {
2651                let mut buffer = Vec::with_capacity(1024);
2652                let mut def_levels = Vec::with_capacity(1024);
2653                let (num_records, num_values, num_levels) = reader
2654                    .read_records(1024, Some(&mut def_levels), None, &mut buffer)
2655                    .unwrap();
2656
2657                assert_eq!(num_records, 513);
2658                assert_eq!(num_values, 513);
2659                assert_eq!(num_levels, 513);
2660
2661                let expected: Vec<i64> = (1..514).collect();
2662                assert_eq!(&buffer, &expected);
2663            }
2664            _ => unreachable!(),
2665        }
2666    }
2667
2668    #[test]
2669    fn test_byte_stream_split_extended() {
2670        let path = format!(
2671            "{}/byte_stream_split_extended.gzip.parquet",
2672            arrow::util::test_util::parquet_test_data(),
2673        );
2674        let file = File::open(path).unwrap();
2675        let reader = Box::new(SerializedFileReader::new(file).expect("Failed to create reader"));
2676
2677        // Use full schema as projected schema
2678        let mut iter = reader
2679            .get_row_iter(None)
2680            .expect("Failed to create row iterator");
2681
2682        let mut start = 0;
2683        let end = reader.metadata().file_metadata().num_rows();
2684
2685        let check_row = |row: Result<Row, ParquetError>| {
2686            assert!(row.is_ok());
2687            let r = row.unwrap();
2688            assert_eq!(r.get_float16(0).unwrap(), r.get_float16(1).unwrap());
2689            assert_eq!(r.get_float(2).unwrap(), r.get_float(3).unwrap());
2690            assert_eq!(r.get_double(4).unwrap(), r.get_double(5).unwrap());
2691            assert_eq!(r.get_int(6).unwrap(), r.get_int(7).unwrap());
2692            assert_eq!(r.get_long(8).unwrap(), r.get_long(9).unwrap());
2693            assert_eq!(r.get_bytes(10).unwrap(), r.get_bytes(11).unwrap());
2694            assert_eq!(r.get_decimal(12).unwrap(), r.get_decimal(13).unwrap());
2695        };
2696
2697        while start < end {
2698            match iter.next() {
2699                Some(row) => check_row(row),
2700                None => break,
2701            };
2702            start += 1;
2703        }
2704    }
2705
2706    #[test]
2707    fn test_filtered_rowgroup_metadata() {
2708        let message_type = "
2709            message test_schema {
2710                REQUIRED INT32 a;
2711            }
2712        ";
2713        let schema = Arc::new(parse_message_type(message_type).unwrap());
2714        let props = Arc::new(
2715            WriterProperties::builder()
2716                .set_statistics_enabled(EnabledStatistics::Page)
2717                .build(),
2718        );
2719        let mut file: File = tempfile::tempfile().unwrap();
2720        let mut file_writer = SerializedFileWriter::new(&mut file, schema, props).unwrap();
2721        let data = [1, 2, 3, 4, 5];
2722
2723        // write 5 row groups
2724        for idx in 0..5 {
2725            let data_i: Vec<i32> = data.iter().map(|x| x * (idx + 1)).collect();
2726            let mut row_group_writer = file_writer.next_row_group().unwrap();
2727            if let Some(mut writer) = row_group_writer.next_column().unwrap() {
2728                writer
2729                    .typed::<Int32Type>()
2730                    .write_batch(data_i.as_slice(), None, None)
2731                    .unwrap();
2732                writer.close().unwrap();
2733            }
2734            row_group_writer.close().unwrap();
2735            file_writer.flushed_row_groups();
2736        }
2737        let file_metadata = file_writer.close().unwrap();
2738
2739        assert_eq!(file_metadata.file_metadata().num_rows(), 25);
2740        assert_eq!(file_metadata.num_row_groups(), 5);
2741
2742        // read only the 3rd row group
2743        let read_options = ReadOptionsBuilder::new()
2744            .with_page_index()
2745            .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) == 2))
2746            .build();
2747        let reader =
2748            SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2749                .unwrap();
2750        let metadata = reader.metadata();
2751
2752        // check we got the expected row group
2753        assert_eq!(metadata.num_row_groups(), 1);
2754        assert_eq!(metadata.row_group(0).ordinal(), Some(2));
2755
2756        // check we only got the relevant page indexes
2757        assert!(metadata.column_index().is_some());
2758        assert!(metadata.offset_index().is_some());
2759        assert_eq!(metadata.column_index().unwrap().len(), 1);
2760        assert_eq!(metadata.offset_index().unwrap().len(), 1);
2761        let col_idx = metadata.column_index().unwrap();
2762        let off_idx = metadata.offset_index().unwrap();
2763        let col_stats = metadata.row_group(0).column(0).statistics().unwrap();
2764        let pg_idx = &col_idx[0][0];
2765        let off_idx_i = &off_idx[0][0];
2766
2767        // test that we got the index matching the row group
2768        match pg_idx {
2769            ColumnIndexMetaData::INT32(int_idx) => {
2770                let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2771                let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2772                assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2773                assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2774            }
2775            _ => panic!("wrong stats type"),
2776        }
2777
2778        // check offset index matches too
2779        assert_eq!(
2780            off_idx_i.page_locations[0].offset,
2781            metadata.row_group(0).column(0).data_page_offset()
2782        );
2783
2784        // read non-contiguous row groups
2785        let read_options = ReadOptionsBuilder::new()
2786            .with_page_index()
2787            .with_predicate(Box::new(|rgmeta, _| rgmeta.ordinal().unwrap_or(0) % 2 == 1))
2788            .build();
2789        let reader =
2790            SerializedFileReader::new_with_options(file.try_clone().unwrap(), read_options)
2791                .unwrap();
2792        let metadata = reader.metadata();
2793
2794        // check we got the expected row groups
2795        assert_eq!(metadata.num_row_groups(), 2);
2796        assert_eq!(metadata.row_group(0).ordinal(), Some(1));
2797        assert_eq!(metadata.row_group(1).ordinal(), Some(3));
2798
2799        // check we only got the relevant page indexes
2800        assert!(metadata.column_index().is_some());
2801        assert!(metadata.offset_index().is_some());
2802        assert_eq!(metadata.column_index().unwrap().len(), 2);
2803        assert_eq!(metadata.offset_index().unwrap().len(), 2);
2804        let col_idx = metadata.column_index().unwrap();
2805        let off_idx = metadata.offset_index().unwrap();
2806
2807        for (i, col_idx_i) in col_idx.iter().enumerate().take(metadata.num_row_groups()) {
2808            let col_stats = metadata.row_group(i).column(0).statistics().unwrap();
2809            let pg_idx = &col_idx_i[0];
2810            let off_idx_i = &off_idx[i][0];
2811
2812            // test that we got the index matching the row group
2813            match pg_idx {
2814                ColumnIndexMetaData::INT32(int_idx) => {
2815                    let min = col_stats.min_bytes_opt().unwrap().get_i32_le();
2816                    let max = col_stats.max_bytes_opt().unwrap().get_i32_le();
2817                    assert_eq!(int_idx.min_value(0), Some(min).as_ref());
2818                    assert_eq!(int_idx.max_value(0), Some(max).as_ref());
2819                }
2820                _ => panic!("wrong stats type"),
2821            }
2822
2823            // check offset index matches too
2824            assert_eq!(
2825                off_idx_i.page_locations[0].offset,
2826                metadata.row_group(i).column(0).data_page_offset()
2827            );
2828        }
2829    }
2830
2831    #[test]
2832    fn test_reuse_schema() {
2833        let file = get_test_file("alltypes_plain.parquet");
2834        let file_reader = SerializedFileReader::new(file.try_clone().unwrap()).unwrap();
2835        let schema = file_reader.metadata().file_metadata().schema_descr_ptr();
2836        let expected = file_reader.metadata;
2837
2838        let options = ReadOptionsBuilder::new()
2839            .with_parquet_schema(schema)
2840            .build();
2841        let file_reader = SerializedFileReader::new_with_options(file, options).unwrap();
2842
2843        assert_eq!(expected.as_ref(), file_reader.metadata.as_ref());
2844        // Should have used the same schema instance
2845        assert!(Arc::ptr_eq(
2846            &expected.file_metadata().schema_descr_ptr(),
2847            &file_reader.metadata.file_metadata().schema_descr_ptr()
2848        ));
2849    }
2850
2851    #[test]
2852    fn test_read_unknown_logical_type() {
2853        let file = get_test_file("unknown-logical-type.parquet");
2854        let reader = SerializedFileReader::new(file).expect("Error opening file");
2855
2856        let schema = reader.metadata().file_metadata().schema_descr();
2857        assert_eq!(
2858            schema.column(0).logical_type_ref(),
2859            Some(&basic::LogicalType::String)
2860        );
2861        assert_eq!(
2862            schema.column(1).logical_type_ref(),
2863            Some(&basic::LogicalType::_Unknown { field_id: 2555 })
2864        );
2865        assert_eq!(schema.column(1).physical_type(), Type::BYTE_ARRAY);
2866
2867        let mut iter = reader
2868            .get_row_iter(None)
2869            .expect("Failed to create row iterator");
2870
2871        let mut num_rows = 0;
2872        while iter.next().is_some() {
2873            num_rows += 1;
2874        }
2875        assert_eq!(num_rows, reader.metadata().file_metadata().num_rows());
2876    }
2877}