Skip to main content

iceberg/arrow/reader/
mod.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18//! Parquet file data reader
19
20use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
21use crate::io::FileIO;
22use crate::runtime::Runtime;
23use crate::util::available_parallelism;
24
25/// Default gap between byte ranges below which they are coalesced into a
26/// single request. Matches object_store's `OBJECT_STORE_COALESCE_DEFAULT`.
27const DEFAULT_RANGE_COALESCE_BYTES: u64 = 1024 * 1024;
28
29/// Default maximum number of coalesced byte ranges fetched concurrently.
30/// Matches object_store's `OBJECT_STORE_COALESCE_PARALLEL`.
31const DEFAULT_RANGE_FETCH_CONCURRENCY: usize = 10;
32
33/// Default number of bytes to prefetch when parsing Parquet footer metadata.
34/// Matches DataFusion's default `ParquetOptions::metadata_size_hint`.
35const DEFAULT_METADATA_SIZE_HINT: usize = 512 * 1024;
36
37mod file_reader;
38mod options;
39mod pipeline;
40mod positional_deletes;
41mod predicate_visitor;
42mod projection;
43mod row_filter;
44mod row_lineage;
45pub use file_reader::ArrowFileReader;
46pub(crate) use options::ParquetReadOptions;
47use predicate_visitor::{CollectFieldIdVisitor, PredicateConverter};
48use projection::{
49    add_fallback_field_ids_to_arrow_schema, apply_name_mapping_to_arrow_schema,
50    find_leaf_by_field_id,
51};
52
53/// Builder to create ArrowReader
54pub struct ArrowReaderBuilder {
55    batch_size: Option<usize>,
56    file_io: FileIO,
57    concurrency_limit_data_files: usize,
58    row_group_filtering_enabled: bool,
59    row_selection_enabled: bool,
60    bloom_filter_enabled: bool,
61    parquet_read_options: ParquetReadOptions,
62    runtime: Runtime,
63}
64
65impl ArrowReaderBuilder {
66    /// Create a new ArrowReaderBuilder
67    pub fn new(file_io: FileIO, runtime: Runtime) -> Self {
68        let num_cpus = available_parallelism().get();
69
70        ArrowReaderBuilder {
71            batch_size: None,
72            file_io,
73            concurrency_limit_data_files: num_cpus,
74            row_group_filtering_enabled: true,
75            row_selection_enabled: false,
76            bloom_filter_enabled: false,
77            parquet_read_options: ParquetReadOptions::builder().build(),
78            runtime,
79        }
80    }
81
82    /// Sets the max number of in flight data files that are being fetched
83    pub fn with_data_file_concurrency_limit(mut self, val: usize) -> Self {
84        self.concurrency_limit_data_files = val;
85        self
86    }
87
88    /// Sets the desired size of batches in the response
89    /// to something other than the default
90    pub fn with_batch_size(mut self, batch_size: usize) -> Self {
91        self.batch_size = Some(batch_size);
92        self
93    }
94
95    /// Determines whether to enable row group filtering.
96    pub fn with_row_group_filtering_enabled(mut self, row_group_filtering_enabled: bool) -> Self {
97        self.row_group_filtering_enabled = row_group_filtering_enabled;
98        self
99    }
100
101    /// Determines whether to enable row selection.
102    pub fn with_row_selection_enabled(mut self, row_selection_enabled: bool) -> Self {
103        self.row_selection_enabled = row_selection_enabled;
104        self
105    }
106
107    /// Determines whether to enable bloom filter-based row group filtering.
108    ///
109    /// When enabled, if a read is performed with an equality or IN predicate,
110    /// the bloom filter for relevant columns in each row group is read and
111    /// checked. Row groups where the bloom filter proves the value is absent
112    /// are skipped entirely.
113    ///
114    /// Defaults to disabled. Each bloom filter is a separate read, and they are
115    /// issued serially — one round trip per relevant column per row group, before
116    /// any data is read. TODO(#3191)
117    pub fn with_bloom_filter_enabled(mut self, bloom_filter_enabled: bool) -> Self {
118        self.bloom_filter_enabled = bloom_filter_enabled;
119        self
120    }
121
122    /// Provide a hint as to the number of bytes to prefetch for parsing the Parquet metadata
123    ///
124    /// This hint can help reduce the number of fetch requests. For more details see the
125    /// [ParquetMetaDataReader documentation](https://docs.rs/parquet/latest/parquet/file/metadata/struct.ParquetMetaDataReader.html#method.with_prefetch_hint).
126    pub fn with_metadata_size_hint(mut self, metadata_size_hint: usize) -> Self {
127        self.parquet_read_options.metadata_size_hint = Some(metadata_size_hint);
128        self
129    }
130
131    /// Sets the gap threshold for merging nearby byte ranges into a single request.
132    /// Ranges with gaps smaller than this value will be coalesced.
133    ///
134    /// Defaults to 1 MiB, matching object_store's OBJECT_STORE_COALESCE_DEFAULT.
135    pub fn with_range_coalesce_bytes(mut self, range_coalesce_bytes: u64) -> Self {
136        self.parquet_read_options.range_coalesce_bytes = range_coalesce_bytes;
137        self
138    }
139
140    /// Sets the maximum number of merged byte ranges to fetch concurrently.
141    ///
142    /// Defaults to 10, matching object_store's OBJECT_STORE_COALESCE_PARALLEL.
143    pub fn with_range_fetch_concurrency(mut self, range_fetch_concurrency: usize) -> Self {
144        self.parquet_read_options.range_fetch_concurrency = range_fetch_concurrency;
145        self
146    }
147
148    /// Build the ArrowReader.
149    pub fn build(self) -> ArrowReader {
150        ArrowReader {
151            batch_size: self.batch_size,
152            file_io: self.file_io.clone(),
153            delete_file_loader: CachingDeleteFileLoader::new(
154                self.file_io.clone(),
155                self.concurrency_limit_data_files,
156                self.runtime.clone(),
157            ),
158            concurrency_limit_data_files: self.concurrency_limit_data_files,
159            row_group_filtering_enabled: self.row_group_filtering_enabled,
160            row_selection_enabled: self.row_selection_enabled,
161            bloom_filter_enabled: self.bloom_filter_enabled,
162            parquet_read_options: self.parquet_read_options,
163        }
164    }
165}
166
167/// Reads data from Parquet files
168#[derive(Clone)]
169pub struct ArrowReader {
170    batch_size: Option<usize>,
171    file_io: FileIO,
172    delete_file_loader: CachingDeleteFileLoader,
173
174    /// the maximum number of data files that can be fetched at the same time
175    concurrency_limit_data_files: usize,
176
177    row_group_filtering_enabled: bool,
178    row_selection_enabled: bool,
179    bloom_filter_enabled: bool,
180    parquet_read_options: ParquetReadOptions,
181}