parquet/arrow/arrow_reader/selection/
ranges.rs1use super::{RowSelection, RowSelector};
28use crate::file::page_index::offset_index::PageLocation;
29use std::ops::Range;
30
31#[inline]
33pub(super) fn scan_ranges_from_selectors<I>(
34 selectors: I,
35 page_locations: &[PageLocation],
36) -> Vec<Range<u64>>
37where
38 I: IntoIterator<Item = RowSelector>,
39{
40 let mut ranges: Vec<Range<u64>> = vec![];
41 let mut row_offset = 0;
42
43 let mut pages = page_locations.iter().peekable();
44 let mut selectors = selectors.into_iter();
45 let mut current_selector = selectors.next();
46 let mut current_page = pages.next();
47
48 let mut current_page_included = false;
49
50 while let Some((selector, page)) = current_selector.as_mut().zip(current_page) {
51 if !(selector.skip || current_page_included) {
52 let start = page.offset as u64;
53 let end = start + page.compressed_page_size as u64;
54 ranges.push(start..end);
55 current_page_included = true;
56 }
57
58 if let Some(next_page) = pages.peek() {
59 if row_offset + selector.row_count > next_page.first_row_index as usize {
60 let remaining_in_page = next_page.first_row_index as usize - row_offset;
61 selector.row_count -= remaining_in_page;
62 row_offset += remaining_in_page;
63 current_page = pages.next();
64 current_page_included = false;
65
66 continue;
67 } else {
68 if row_offset + selector.row_count == next_page.first_row_index as usize {
69 current_page = pages.next();
70 current_page_included = false;
71 }
72 row_offset += selector.row_count;
73 current_selector = selectors.next();
74 }
75 } else {
76 if !(selector.skip || current_page_included) {
77 let start = page.offset as u64;
78 let end = start + page.compressed_page_size as u64;
79 ranges.push(start..end);
80 }
81 current_selector = selectors.next()
82 }
83 }
84
85 ranges
86}
87
88#[inline]
91pub(super) fn expand_to_batch_boundaries_from_selectors<I>(
92 selectors: I,
93 batch_size: usize,
94 total_rows: usize,
95) -> RowSelection
96where
97 I: IntoIterator<Item = RowSelector>,
98{
99 let mut expanded_ranges = Vec::new();
100 let mut row_offset = 0;
101
102 for selector in selectors {
103 if selector.skip {
104 row_offset += selector.row_count;
105 } else {
106 let start = row_offset;
107 let end = row_offset + selector.row_count;
108
109 let expanded_start = (start / batch_size) * batch_size;
111 let expanded_end = end.div_ceil(batch_size) * batch_size;
113 let expanded_end = expanded_end.min(total_rows);
114
115 expanded_ranges.push(expanded_start..expanded_end);
116 row_offset += selector.row_count;
117 }
118 }
119
120 expanded_ranges.sort_by_key(|range| range.start);
122
123 let mut merged_ranges: Vec<Range<usize>> = Vec::new();
125 for range in expanded_ranges {
126 if let Some(last) = merged_ranges.last_mut() {
127 if range.start <= last.end {
128 last.end = last.end.max(range.end);
130 } else {
131 merged_ranges.push(range);
133 }
134 } else {
135 merged_ranges.push(range);
137 }
138 }
139
140 RowSelection::from_consecutive_ranges(merged_ranges.into_iter(), total_rows)
141}
142
143#[cfg(test)]
144mod tests {
145 use super::*;
146
147 #[test]
148 fn test_scan_ranges() {
149 let index = vec![
150 PageLocation {
151 offset: 0,
152 compressed_page_size: 10,
153 first_row_index: 0,
154 },
155 PageLocation {
156 offset: 10,
157 compressed_page_size: 10,
158 first_row_index: 10,
159 },
160 PageLocation {
161 offset: 20,
162 compressed_page_size: 10,
163 first_row_index: 20,
164 },
165 PageLocation {
166 offset: 30,
167 compressed_page_size: 10,
168 first_row_index: 30,
169 },
170 PageLocation {
171 offset: 40,
172 compressed_page_size: 10,
173 first_row_index: 40,
174 },
175 PageLocation {
176 offset: 50,
177 compressed_page_size: 10,
178 first_row_index: 50,
179 },
180 PageLocation {
181 offset: 60,
182 compressed_page_size: 10,
183 first_row_index: 60,
184 },
185 ];
186
187 let selection = RowSelection::from(vec![
188 RowSelector::skip(10),
190 RowSelector::select(3),
192 RowSelector::skip(3),
193 RowSelector::select(4),
194 RowSelector::skip(5),
196 RowSelector::select(5),
197 RowSelector::skip(12),
199 RowSelector::select(12),
201 RowSelector::skip(12),
203 ]);
204
205 let ranges = selection.scan_ranges(&index);
206
207 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60]);
209 assert_eq!(
210 selection.row_ranges_for_selected_pages(&index, 70),
211 vec![10..20, 20..30, 40..50, 50..60]
212 );
213
214 let selection = RowSelection::from(vec![
215 RowSelector::skip(10),
217 RowSelector::select(3),
219 RowSelector::skip(3),
220 RowSelector::select(4),
221 RowSelector::skip(5),
223 RowSelector::select(5),
224 RowSelector::skip(12),
226 RowSelector::select(12),
228 RowSelector::skip(1),
229 RowSelector::select(8),
231 ]);
232
233 let ranges = selection.scan_ranges(&index);
234
235 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60, 60..70]);
237
238 let selection = RowSelection::from(vec![
239 RowSelector::skip(10),
241 RowSelector::select(3),
243 RowSelector::skip(3),
244 RowSelector::select(4),
245 RowSelector::skip(5),
247 RowSelector::select(5),
248 RowSelector::skip(12),
250 RowSelector::select(12),
252 RowSelector::skip(1),
253 RowSelector::skip(8),
255 RowSelector::select(4),
257 ]);
258
259 let ranges = selection.scan_ranges(&index);
260
261 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60, 60..70]);
263
264 let selection = RowSelection::from(vec![
265 RowSelector::skip(10),
267 RowSelector::select(3),
269 RowSelector::skip(3),
270 RowSelector::select(4),
271 RowSelector::skip(5),
273 RowSelector::select(6),
274 RowSelector::skip(50),
276 ]);
277
278 let ranges = selection.scan_ranges(&index);
279
280 assert_eq!(ranges, vec![10..20, 20..30, 30..40]);
282 }
283}