iceberg/arrow/reader/
mod.rs1use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
21use crate::io::FileIO;
22use crate::runtime::Runtime;
23use crate::util::available_parallelism;
24
25const DEFAULT_RANGE_COALESCE_BYTES: u64 = 1024 * 1024;
28
29const DEFAULT_RANGE_FETCH_CONCURRENCY: usize = 10;
32
33const 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
53pub 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 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 pub fn with_data_file_concurrency_limit(mut self, val: usize) -> Self {
84 self.concurrency_limit_data_files = val;
85 self
86 }
87
88 pub fn with_batch_size(mut self, batch_size: usize) -> Self {
91 self.batch_size = Some(batch_size);
92 self
93 }
94
95 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 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 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 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 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 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 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#[derive(Clone)]
169pub struct ArrowReader {
170 batch_size: Option<usize>,
171 file_io: FileIO,
172 delete_file_loader: CachingDeleteFileLoader,
173
174 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}