parquet/arrow/push_decoder/reader_builder/
data.rs1use crate::arrow::ProjectionMask;
21use crate::arrow::arrow_reader::RowSelection;
22use crate::arrow::in_memory_row_group::{ColumnChunkData, FetchRanges, InMemoryRowGroup};
23use crate::errors::ParquetError;
24use crate::file::metadata::ParquetMetaData;
25use crate::file::reader::ChunkReader;
26use crate::util::push_buffers::PushBuffers;
27use bytes::Bytes;
28use std::ops::Range;
29use std::sync::Arc;
30
31#[derive(Debug)]
35pub(super) struct DataRequest {
36 column_chunks: Vec<Option<Arc<ColumnChunkData>>>,
38 ranges: Vec<Range<u64>>,
40 page_start_offsets: Option<Vec<Vec<u64>>>,
43}
44
45impl DataRequest {
46 pub fn needed_ranges(&self, buffers: &PushBuffers) -> Vec<Range<u64>> {
49 self.ranges
50 .iter()
51 .filter(|&range| !buffers.has_range(range))
52 .cloned()
53 .collect()
54 }
55
56 fn get_chunks(&self, buffers: &PushBuffers) -> Result<Vec<Bytes>, ParquetError> {
58 self.ranges
59 .iter()
60 .map(|range| {
61 let length: usize = (range.end - range.start)
62 .try_into()
63 .expect("overflow for offset");
64 buffers.get_bytes(range.start, length).map_err(|e| {
66 ParquetError::General(format!(
67 "Internal Error missing data for range {range:?} in buffers: {e}",
68 ))
69 })
70 })
71 .collect()
72 }
73
74 pub fn try_into_in_memory_row_group<'a>(
79 self,
80 row_group_idx: usize,
81 row_count: usize,
82 parquet_metadata: &'a ParquetMetaData,
83 projection: &ProjectionMask,
84 buffers: &mut PushBuffers,
85 ) -> Result<InMemoryRowGroup<'a>, ParquetError> {
86 let chunks = self.get_chunks(buffers)?;
87
88 let Self {
89 column_chunks,
90 ranges,
91 page_start_offsets,
92 } = self;
93
94 let page_index = if parquet_metadata
95 .page_index()
96 .is_some_and(|pi| pi.has_offset_indexes())
97 {
98 Some(parquet_metadata.page_index_for_row_group(row_group_idx))
99 } else {
100 None
101 };
102
103 let mut in_memory_row_group = InMemoryRowGroup {
107 row_count,
108 column_chunks,
109 page_index,
110 row_group_idx,
111 metadata: parquet_metadata,
112 };
113
114 in_memory_row_group.fill_column_chunks(projection, page_start_offsets, chunks);
115
116 buffers.clear_ranges(&ranges);
118
119 Ok(in_memory_row_group)
120 }
121}
122
123pub(super) struct DataRequestBuilder<'a> {
125 row_group_idx: usize,
127 row_count: usize,
129 batch_size: usize,
131 parquet_metadata: &'a ParquetMetaData,
133 projection: &'a ProjectionMask,
135 selection: Option<&'a RowSelection>,
137 cache_projection: Option<&'a ProjectionMask>,
141 column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
143}
144
145impl<'a> DataRequestBuilder<'a> {
146 pub(super) fn new(
147 row_group_idx: usize,
148 row_count: usize,
149 batch_size: usize,
150 parquet_metadata: &'a ParquetMetaData,
151 projection: &'a ProjectionMask,
152 ) -> Self {
153 Self {
154 row_group_idx,
155 row_count,
156 batch_size,
157 parquet_metadata,
158 projection,
159 selection: None,
160 cache_projection: None,
161 column_chunks: None,
162 }
163 }
164
165 pub(super) fn with_selection(mut self, selection: Option<&'a RowSelection>) -> Self {
167 self.selection = selection;
168 self
169 }
170
171 pub(super) fn with_cache_projection(
173 mut self,
174 cache_projection: Option<&'a ProjectionMask>,
175 ) -> Self {
176 self.cache_projection = cache_projection;
177 self
178 }
179
180 pub(super) fn with_column_chunks(
182 mut self,
183 column_chunks: Option<Vec<Option<Arc<ColumnChunkData>>>>,
184 ) -> Self {
185 self.column_chunks = column_chunks;
186 self
187 }
188
189 pub(crate) fn build(self) -> DataRequest {
190 let Self {
191 row_group_idx,
192 row_count,
193 batch_size,
194 parquet_metadata,
195 projection,
196 selection,
197 cache_projection,
198 column_chunks,
199 } = self;
200
201 let row_group_meta_data = parquet_metadata.row_group(row_group_idx);
202
203 let column_chunks =
205 column_chunks.unwrap_or_else(|| vec![None; row_group_meta_data.columns().len()]);
206
207 let page_index = if parquet_metadata
208 .page_index()
209 .is_some_and(|pi| pi.has_offset_indexes())
210 {
211 Some(parquet_metadata.page_index_for_row_group(row_group_idx))
212 } else {
213 None
214 };
215
216 let row_group = InMemoryRowGroup {
220 row_count,
221 column_chunks,
222 page_index,
223 row_group_idx,
224 metadata: parquet_metadata,
225 };
226
227 let FetchRanges {
228 ranges,
229 page_start_offsets,
230 } = row_group.fetch_ranges(projection, selection, batch_size, cache_projection);
231
232 DataRequest {
233 column_chunks: row_group.column_chunks,
235 ranges,
236 page_start_offsets,
237 }
238 }
239}