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