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 } else {
66 if row_offset + selector.row_count == next_page.first_row_index as usize {
67 current_page = pages.next();
68 current_page_included = false;
69 }
70 row_offset += selector.row_count;
71 current_selector = selectors.next();
72 }
73 } else {
74 if !(selector.skip || current_page_included) {
75 let start = page.offset as u64;
76 let end = start + page.compressed_page_size as u64;
77 ranges.push(start..end);
78 }
79 current_selector = selectors.next()
80 }
81 }
82
83 ranges
84}
85
86#[inline]
89pub(super) fn expand_to_batch_boundaries_from_selectors<I>(
90 selectors: I,
91 batch_size: usize,
92 total_rows: usize,
93) -> RowSelection
94where
95 I: IntoIterator<Item = RowSelector>,
96{
97 let mut expanded_ranges = Vec::new();
98 let mut row_offset = 0;
99
100 for selector in selectors {
101 if !selector.skip {
102 let start = row_offset;
103 let end = row_offset + selector.row_count;
104
105 let expanded_start = (start / batch_size) * batch_size;
107 let expanded_end = end.div_ceil(batch_size) * batch_size;
109 let expanded_end = expanded_end.min(total_rows);
110
111 expanded_ranges.push(expanded_start..expanded_end);
112 }
113 row_offset += selector.row_count;
114 }
115
116 expanded_ranges.sort_by_key(|range| range.start);
118
119 let mut merged_ranges: Vec<Range<usize>> = Vec::new();
121 for range in expanded_ranges {
122 if let Some(last) = merged_ranges.last_mut() {
123 if range.start <= last.end {
124 last.end = last.end.max(range.end);
126 } else {
127 merged_ranges.push(range);
129 }
130 } else {
131 merged_ranges.push(range);
133 }
134 }
135
136 RowSelection::from_consecutive_ranges(merged_ranges.into_iter(), total_rows)
137}
138
139#[cfg(test)]
140mod tests {
141 use super::*;
142
143 #[test]
144 fn test_scan_ranges() {
145 let index = vec![
146 PageLocation {
147 offset: 0,
148 compressed_page_size: 10,
149 first_row_index: 0,
150 },
151 PageLocation {
152 offset: 10,
153 compressed_page_size: 10,
154 first_row_index: 10,
155 },
156 PageLocation {
157 offset: 20,
158 compressed_page_size: 10,
159 first_row_index: 20,
160 },
161 PageLocation {
162 offset: 30,
163 compressed_page_size: 10,
164 first_row_index: 30,
165 },
166 PageLocation {
167 offset: 40,
168 compressed_page_size: 10,
169 first_row_index: 40,
170 },
171 PageLocation {
172 offset: 50,
173 compressed_page_size: 10,
174 first_row_index: 50,
175 },
176 PageLocation {
177 offset: 60,
178 compressed_page_size: 10,
179 first_row_index: 60,
180 },
181 ];
182
183 let selection = RowSelection::from(vec![
184 RowSelector::skip(10),
186 RowSelector::select(3),
188 RowSelector::skip(3),
189 RowSelector::select(4),
190 RowSelector::skip(5),
192 RowSelector::select(5),
193 RowSelector::skip(12),
195 RowSelector::select(12),
197 RowSelector::skip(12),
199 ]);
200
201 let ranges = selection.scan_ranges(&index);
202
203 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60]);
205 assert_eq!(
206 selection.row_ranges_for_selected_pages(&index, 70),
207 vec![10..20, 20..30, 40..50, 50..60]
208 );
209
210 let selection = RowSelection::from(vec![
211 RowSelector::skip(10),
213 RowSelector::select(3),
215 RowSelector::skip(3),
216 RowSelector::select(4),
217 RowSelector::skip(5),
219 RowSelector::select(5),
220 RowSelector::skip(12),
222 RowSelector::select(12),
224 RowSelector::skip(1),
225 RowSelector::select(8),
227 ]);
228
229 let ranges = selection.scan_ranges(&index);
230
231 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60, 60..70]);
233
234 let selection = RowSelection::from(vec![
235 RowSelector::skip(10),
237 RowSelector::select(3),
239 RowSelector::skip(3),
240 RowSelector::select(4),
241 RowSelector::skip(5),
243 RowSelector::select(5),
244 RowSelector::skip(12),
246 RowSelector::select(12),
248 RowSelector::skip(1),
249 RowSelector::skip(8),
251 RowSelector::select(4),
253 ]);
254
255 let ranges = selection.scan_ranges(&index);
256
257 assert_eq!(ranges, vec![10..20, 20..30, 40..50, 50..60, 60..70]);
259
260 let selection = RowSelection::from(vec![
261 RowSelector::skip(10),
263 RowSelector::select(3),
265 RowSelector::skip(3),
266 RowSelector::select(4),
267 RowSelector::skip(5),
269 RowSelector::select(6),
270 RowSelector::skip(50),
272 ]);
273
274 let ranges = selection.scan_ranges(&index);
275
276 assert_eq!(ranges, vec![10..20, 20..30, 30..40]);
278 }
279}