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 parquet_read_options: ParquetReadOptions,
61 runtime: Runtime,
62}
63
64impl ArrowReaderBuilder {
65 pub fn new(file_io: FileIO, runtime: Runtime) -> Self {
67 let num_cpus = available_parallelism().get();
68
69 ArrowReaderBuilder {
70 batch_size: None,
71 file_io,
72 concurrency_limit_data_files: num_cpus,
73 row_group_filtering_enabled: true,
74 row_selection_enabled: false,
75 parquet_read_options: ParquetReadOptions::builder().build(),
76 runtime,
77 }
78 }
79
80 pub fn with_data_file_concurrency_limit(mut self, val: usize) -> Self {
82 self.concurrency_limit_data_files = val;
83 self
84 }
85
86 pub fn with_batch_size(mut self, batch_size: usize) -> Self {
89 self.batch_size = Some(batch_size);
90 self
91 }
92
93 pub fn with_row_group_filtering_enabled(mut self, row_group_filtering_enabled: bool) -> Self {
95 self.row_group_filtering_enabled = row_group_filtering_enabled;
96 self
97 }
98
99 pub fn with_row_selection_enabled(mut self, row_selection_enabled: bool) -> Self {
101 self.row_selection_enabled = row_selection_enabled;
102 self
103 }
104
105 pub fn with_metadata_size_hint(mut self, metadata_size_hint: usize) -> Self {
110 self.parquet_read_options.metadata_size_hint = Some(metadata_size_hint);
111 self
112 }
113
114 pub fn with_range_coalesce_bytes(mut self, range_coalesce_bytes: u64) -> Self {
119 self.parquet_read_options.range_coalesce_bytes = range_coalesce_bytes;
120 self
121 }
122
123 pub fn with_range_fetch_concurrency(mut self, range_fetch_concurrency: usize) -> Self {
127 self.parquet_read_options.range_fetch_concurrency = range_fetch_concurrency;
128 self
129 }
130
131 pub fn build(self) -> ArrowReader {
133 ArrowReader {
134 batch_size: self.batch_size,
135 file_io: self.file_io.clone(),
136 delete_file_loader: CachingDeleteFileLoader::new(
137 self.file_io.clone(),
138 self.concurrency_limit_data_files,
139 self.runtime.clone(),
140 ),
141 concurrency_limit_data_files: self.concurrency_limit_data_files,
142 row_group_filtering_enabled: self.row_group_filtering_enabled,
143 row_selection_enabled: self.row_selection_enabled,
144 parquet_read_options: self.parquet_read_options,
145 }
146 }
147}
148
149#[derive(Clone)]
151pub struct ArrowReader {
152 batch_size: Option<usize>,
153 file_io: FileIO,
154 delete_file_loader: CachingDeleteFileLoader,
155
156 concurrency_limit_data_files: usize,
158
159 row_group_filtering_enabled: bool,
160 row_selection_enabled: bool,
161 parquet_read_options: ParquetReadOptions,
162}