1use std::borrow::Cow;
2
3use arrow::bitmap::Bitmap;
4use arrow::datatypes::ArrowSchemaRef;
5use polars_buffer::Buffer;
6use polars_core::chunked_array::builder::NullChunkedBuilder;
7use polars_core::config;
8use polars_core::prelude::*;
9use polars_core::runtime::RAYON;
10use polars_core::series::IsSorted;
11use polars_core::utils::accumulate_dataframes_vertical;
12use polars_parquet::read::{self, ColumnChunkMetadata, FileMetadata, Filter, RowGroupMetadata};
13use rayon::prelude::*;
14
15use super::mmap::mmap_columns;
16use super::utils::{canonicalize_parquet_maps, materialize_empty_df};
17use super::{ParallelStrategy, mmap};
18use crate::RowIndex;
19use crate::hive::materialize_hive_partitions;
20use crate::mmap::{MmapBytesReader, ReaderBytes};
21use crate::parquet::metadata::FileMetadataRef;
22use crate::parquet::read::ROW_COUNT_OVERFLOW_ERR;
23use crate::utils::slice::split_slice_at_file;
24
25#[cfg(debug_assertions)]
26fn assert_dtypes(dtype: &ArrowDataType) {
29 use ArrowDataType as D;
30
31 match dtype {
32 D::Utf8 | D::Binary | D::LargeUtf8 | D::LargeBinary => unreachable!(),
34
35 D::Float16 => unreachable!(),
37
38 D::List(_) => unreachable!(),
40
41 D::Map(entries, _) => assert_dtypes(&entries.dtype),
43 D::Dictionary(_, dtype, _) => assert_dtypes(dtype),
44 D::Extension(ext) => assert_dtypes(&ext.inner),
45 D::LargeList(inner) => assert_dtypes(&inner.dtype),
46 D::FixedSizeList(inner, _) => assert_dtypes(&inner.dtype),
47 D::Struct(fields) => fields.iter().for_each(|f| assert_dtypes(f.dtype())),
48
49 _ => {},
50 }
51}
52
53fn should_copy_sortedness(dtype: &DataType) -> bool {
54 use DataType as D;
56
57 matches!(
58 dtype,
59 D::Int8 | D::Int16 | D::Int32 | D::Int64 | D::UInt8 | D::UInt16 | D::UInt32 | D::UInt64
60 )
61}
62
63pub fn try_set_sorted_flag(series: &mut Series, col_idx: usize, sorting_map: &[(usize, IsSorted)]) {
64 let Some((sorted_col, is_sorted)) = sorting_map.first() else {
65 return;
66 };
67 if *sorted_col != col_idx || !should_copy_sortedness(series.dtype()) {
68 return;
69 }
70 if config::verbose() {
71 eprintln!(
72 "Parquet conserved SortingColumn for column chunk of '{}' to {is_sorted:?}",
73 series.name()
74 );
75 }
76
77 series.set_sorted_flag(*is_sorted);
78}
79
80pub fn create_sorting_map(md: &RowGroupMetadata) -> Vec<(usize, IsSorted)> {
81 let capacity = md.sorting_columns().map_or(0, |s| s.len());
82 let mut sorting_map = Vec::with_capacity(capacity);
83
84 if let Some(sorting_columns) = md.sorting_columns() {
85 for sorting in sorting_columns {
86 sorting_map.push((
87 sorting.column_idx as usize,
88 if sorting.descending {
89 IsSorted::Descending
90 } else {
91 IsSorted::Ascending
92 },
93 ))
94 }
95 }
96
97 sorting_map
98}
99
100fn column_idx_to_series(
101 column_i: usize,
102 field_md: &[&ColumnChunkMetadata],
104 filter: Option<Filter>,
105 file_schema: &ArrowSchema,
106 store: &mmap::ColumnStore,
107) -> PolarsResult<(Series, Bitmap)> {
108 let field = file_schema.get_at_index(column_i).unwrap().1;
109
110 #[cfg(debug_assertions)]
111 {
112 assert_dtypes(field.dtype())
113 }
114 let columns = mmap_columns(store, field_md);
115 let (arrays, pred_true_mask) = mmap::to_deserializer(columns, field.clone(), filter)?;
116 let mut series = Series::try_from((field, arrays))?;
117 canonicalize_parquet_maps(&mut series)?;
118
119 Ok((series, pred_true_mask))
120}
121
122#[allow(clippy::too_many_arguments)]
123fn rg_to_dfs(
124 store: &mmap::ColumnStore,
125 previous_row_count: &mut IdxSize,
126 row_group_start: usize,
127 row_group_end: usize,
128 pre_slice: (usize, usize),
129 file_metadata: &FileMetadata,
130 schema: &ArrowSchemaRef,
131 row_index: Option<RowIndex>,
132 parallel: ParallelStrategy,
133 projection: &[usize],
134 hive_partition_columns: Option<&[Series]>,
135) -> PolarsResult<Vec<DataFrame>> {
136 if config::verbose() {
137 eprintln!("parquet scan with parallel = {parallel:?}");
138 }
139
140 if projection.is_empty() {
142 if let Some(row_index) = row_index {
143 let placeholder =
144 NullChunkedBuilder::new(PlSmallStr::from_static("__PL_TMP"), pre_slice.1).finish();
145 return Ok(vec![
146 DataFrame::new_infer_height(vec![placeholder.into_series().into_column()])?
147 .with_row_index(
148 row_index.name.clone(),
149 Some(row_index.offset + IdxSize::try_from(pre_slice.0).unwrap()),
150 )?
151 .select(std::iter::once(row_index.name))?,
152 ]);
153 }
154 }
155
156 use ParallelStrategy as S;
157
158 match parallel {
159 S::Columns | S::None => rg_to_dfs_optionally_par_over_columns(
160 store,
161 previous_row_count,
162 row_group_start,
163 row_group_end,
164 pre_slice,
165 file_metadata,
166 schema,
167 row_index,
168 parallel,
169 projection,
170 hive_partition_columns,
171 ),
172 _ => rg_to_dfs_par_over_rg(
173 store,
174 row_group_start,
175 row_group_end,
176 previous_row_count,
177 pre_slice,
178 file_metadata,
179 schema,
180 row_index,
181 projection,
182 hive_partition_columns,
183 ),
184 }
185}
186
187#[allow(clippy::too_many_arguments)]
188fn rg_to_dfs_optionally_par_over_columns(
190 store: &mmap::ColumnStore,
191 previous_row_count: &mut IdxSize,
192 row_group_start: usize,
193 row_group_end: usize,
194 slice: (usize, usize),
195 file_metadata: &FileMetadata,
196 schema: &ArrowSchemaRef,
197 row_index: Option<RowIndex>,
198 parallel: ParallelStrategy,
199 projection: &[usize],
200 hive_partition_columns: Option<&[Series]>,
201) -> PolarsResult<Vec<DataFrame>> {
202 let mut dfs = Vec::with_capacity(row_group_end - row_group_start);
203
204 let mut n_rows_processed: usize = (0..row_group_start)
205 .map(|i| file_metadata.row_groups[i].num_rows())
206 .sum();
207 let slice_end = slice.0 + slice.1;
208
209 for md in &file_metadata.row_groups[row_group_start..row_group_end] {
210 let rg_slice =
211 split_slice_at_file(&mut n_rows_processed, md.num_rows(), slice.0, slice_end);
212 let current_row_count = md.num_rows() as IdxSize;
213
214 let sorting_map = create_sorting_map(md);
215
216 let f = |column_i: &usize| {
217 let (name, field) = schema.get_at_index(*column_i).unwrap();
218
219 let Some(iter) = md.columns_under_root_iter(name) else {
220 return Ok(Column::full_null(
221 name.clone(),
222 rg_slice.1,
223 &DataType::from_arrow_field(field),
224 ));
225 };
226
227 let part = iter.collect::<Vec<_>>();
228
229 let (mut series, _) = column_idx_to_series(
230 *column_i,
231 part.as_slice(),
232 Some(Filter::new_ranged(rg_slice.0, rg_slice.0 + rg_slice.1)),
233 schema,
234 store,
235 )?;
236
237 try_set_sorted_flag(&mut series, *column_i, &sorting_map);
238 Ok(series.into_column())
239 };
240
241 let columns = if let ParallelStrategy::Columns = parallel {
242 RAYON.install(|| {
243 projection
244 .par_iter()
245 .map(f)
246 .collect::<PolarsResult<Vec<_>>>()
247 })?
248 } else {
249 projection.iter().map(f).collect::<PolarsResult<Vec<_>>>()?
250 };
251
252 let mut df = unsafe { DataFrame::new_unchecked(rg_slice.1, columns) };
253 if let Some(rc) = &row_index {
254 unsafe {
255 df.with_row_index_mut(
256 rc.name.clone(),
257 Some(*previous_row_count + rc.offset + rg_slice.0 as IdxSize),
258 )
259 };
260 }
261
262 materialize_hive_partitions(&mut df, schema.as_ref(), hive_partition_columns);
263
264 *previous_row_count = previous_row_count
265 .checked_add(current_row_count)
266 .ok_or_else(|| {
267 polars_err!(
268 ComputeError: "Parquet file produces more than pow(2, 32) rows; \
269 consider compiling with polars-bigidx feature (pip install polars[rt64]), \
270 or set 'streaming'"
271 )
272 })?;
273 dfs.push(df);
274
275 if *previous_row_count as usize >= slice_end {
276 break;
277 }
278 }
279
280 Ok(dfs)
281}
282
283#[allow(clippy::too_many_arguments)]
284fn rg_to_dfs_par_over_rg(
286 store: &mmap::ColumnStore,
287 row_group_start: usize,
288 row_group_end: usize,
289 rows_read: &mut IdxSize,
290 slice: (usize, usize),
291 file_metadata: &FileMetadata,
292 schema: &ArrowSchemaRef,
293 row_index: Option<RowIndex>,
294 projection: &[usize],
295 hive_partition_columns: Option<&[Series]>,
296) -> PolarsResult<Vec<DataFrame>> {
297 let mut row_groups = Vec::with_capacity(row_group_end - row_group_start);
299
300 let mut n_rows_processed: usize = (0..row_group_start)
301 .map(|i| file_metadata.row_groups[i].num_rows())
302 .sum();
303 let slice_end = slice.0 + slice.1;
304
305 let mut rows_scanned: IdxSize;
309
310 if row_group_start > 0 {
311 rows_scanned = (0..row_group_start)
315 .map(|i| file_metadata.row_groups[i].num_rows() as IdxSize)
316 .sum();
317 } else {
318 rows_scanned = 0;
319 }
320
321 for rg_md in &file_metadata.row_groups[row_group_start..row_group_end] {
322 let row_count_start = rows_scanned;
323 let n_rows_this_file = rg_md.num_rows();
324 let rg_slice =
325 split_slice_at_file(&mut n_rows_processed, n_rows_this_file, slice.0, slice_end);
326 rows_scanned = rows_scanned
327 .checked_add(n_rows_this_file as IdxSize)
328 .ok_or(ROW_COUNT_OVERFLOW_ERR)?;
329
330 *rows_read += rg_slice.1 as IdxSize;
331
332 if rg_slice.1 == 0 {
333 continue;
334 }
335
336 row_groups.push((rg_md, rg_slice, row_count_start));
337 }
338
339 let dfs = RAYON.install(|| {
340 row_groups
343 .into_par_iter()
344 .map(|(md, slice, row_count_start)| {
345 if slice.1 == 0 {
346 return Ok(None);
347 }
348 #[cfg(debug_assertions)]
350 {
351 assert!(std::env::var("POLARS_PANIC_IF_PARQUET_PARSED").is_err())
352 }
353
354 let sorting_map = create_sorting_map(md);
355
356 let columns = projection
357 .iter()
358 .map(|column_i| {
359 let (name, field) = schema.get_at_index(*column_i).unwrap();
360
361 let Some(iter) = md.columns_under_root_iter(name) else {
362 return Ok(Column::full_null(
363 name.clone(),
364 md.num_rows(),
365 &DataType::from_arrow_field(field),
366 ));
367 };
368
369 let part = iter.collect::<Vec<_>>();
370
371 let (mut series, _) = column_idx_to_series(
372 *column_i,
373 part.as_slice(),
374 Some(Filter::new_ranged(slice.0, slice.0 + slice.1)),
375 schema,
376 store,
377 )?;
378
379 try_set_sorted_flag(&mut series, *column_i, &sorting_map);
380 Ok(series.into_column())
381 })
382 .collect::<PolarsResult<Vec<_>>>()?;
383
384 let mut df = unsafe { DataFrame::new_unchecked(slice.1, columns) };
385
386 if let Some(rc) = &row_index {
387 unsafe {
388 df.with_row_index_mut(
389 rc.name.clone(),
390 Some(row_count_start as IdxSize + rc.offset + slice.0 as IdxSize),
391 )
392 };
393 }
394
395 materialize_hive_partitions(&mut df, schema.as_ref(), hive_partition_columns);
396
397 Ok(Some(df))
398 })
399 .collect::<PolarsResult<Vec<_>>>()
400 })?;
401 Ok(dfs.into_iter().flatten().collect())
402}
403
404#[allow(clippy::too_many_arguments)]
405pub fn read_parquet<R: MmapBytesReader>(
406 mut reader: R,
407 pre_slice: (usize, usize),
408 projection: Option<&[usize]>,
409 reader_schema: &ArrowSchemaRef,
410 metadata: Option<FileMetadataRef>,
411 mut parallel: ParallelStrategy,
412 row_index: Option<RowIndex>,
413 hive_partition_columns: Option<&[Series]>,
414) -> PolarsResult<DataFrame> {
415 if pre_slice.1 == 0 {
417 return Ok(materialize_empty_df(
418 projection,
419 reader_schema,
420 hive_partition_columns,
421 row_index.as_ref(),
422 ));
423 }
424
425 let file_metadata = metadata
426 .map(Ok)
427 .unwrap_or_else(|| read::read_metadata(&mut reader).map(Arc::new))?;
428 let n_row_groups = file_metadata.row_groups.len();
429
430 let materialized_projection = projection
431 .map(Cow::Borrowed)
432 .unwrap_or_else(|| Cow::Owned((0usize..reader_schema.len()).collect::<Vec<_>>()));
433
434 if ParallelStrategy::Auto == parallel {
435 if n_row_groups > materialized_projection.len()
436 || n_row_groups > RAYON.current_num_threads()
437 {
438 parallel = ParallelStrategy::RowGroups;
439 } else {
440 parallel = ParallelStrategy::Columns;
441 }
442 }
443
444 if let (ParallelStrategy::Columns, true) = (parallel, materialized_projection.len() == 1) {
445 parallel = ParallelStrategy::None;
446 }
447
448 let reader = ReaderBytes::from(&mut reader);
449 Buffer::with_slice(&reader, |buf| {
450 let store = mmap::ColumnStore::Local(buf);
451 let dfs = rg_to_dfs(
452 &store,
453 &mut 0,
454 0,
455 n_row_groups,
456 pre_slice,
457 &file_metadata,
458 reader_schema,
459 row_index.clone(),
460 parallel,
461 &materialized_projection,
462 hive_partition_columns,
463 )?;
464
465 if dfs.is_empty() {
466 Ok(materialize_empty_df(
467 projection,
468 reader_schema,
469 hive_partition_columns,
470 row_index.as_ref(),
471 ))
472 } else {
473 accumulate_dataframes_vertical(dfs)
474 }
475 })
476}
477
478pub fn calc_prefilter_cost(mask: &arrow::bitmap::Bitmap) -> f64 {
479 let num_edges = mask.num_edges() as f64;
480 let rg_len = mask.len() as f64;
481
482 (num_edges / rg_len).clamp(0.0, 1.0)
499}
500
501#[derive(Clone, Copy)]
502pub enum PrefilterMaskSetting {
503 Auto,
504 Pre,
505 Post,
506}
507
508impl PrefilterMaskSetting {
509 pub fn init_from_env() -> Self {
510 std::env::var("POLARS_PQ_PREFILTERED_MASK").map_or(Self::Auto, |v| match &v[..] {
511 "auto" => Self::Auto,
512 "pre" => Self::Pre,
513 "post" => Self::Post,
514 _ => panic!("Invalid `POLARS_PQ_PREFILTERED_MASK` value '{v}'."),
515 })
516 }
517
518 pub fn should_prefilter(&self, prefilter_cost: f64, dtype: &ArrowDataType) -> bool {
519 match self {
520 Self::Auto => {
521 let is_nested = dtype.is_nested();
524
525 !is_nested && prefilter_cost <= 0.01
527 },
528 Self::Pre => true,
529 Self::Post => false,
530 }
531 }
532}