Skip to main content

polars_io/parquet/read/
read_impl.rs

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)]
26// Ensure we get the proper polars types from schema inference
27// This saves unneeded casts.
28fn assert_dtypes(dtype: &ArrowDataType) {
29    use ArrowDataType as D;
30
31    match dtype {
32        // These should all be cast to the BinaryView / Utf8View variants
33        D::Utf8 | D::Binary | D::LargeUtf8 | D::LargeBinary => unreachable!(),
34
35        // These should be cast to Float32
36        D::Float16 => unreachable!(),
37
38        // This should have been converted to a LargeList
39        D::List(_) => unreachable!(),
40
41        // Recursive checks
42        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    // @NOTE: For now, we are a bit conservative with this.
55    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    // The metadata belonging to this column
103    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 we are only interested in the row_index, we take a little special path here.
141    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)]
188// might parallelize over columns
189fn 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)]
284// parallelizes over row groups
285fn 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    // compute the limits per row group and the row count offsets
298    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    // rows_scanned is the number of rows that have been scanned so far when checking for overlap with the slice.
306    // rows_read is the number of rows found to overlap with the slice, and thus the number of rows that will be
307    // read into a dataframe.
308    let mut rows_scanned: IdxSize;
309
310    if row_group_start > 0 {
311        // In the case of async reads, we need to account for the fact that row_group_start may be greater than
312        // zero due to earlier processing.
313        // For details, see: https://github.com/pola-rs/polars/pull/20508#discussion_r1900165649
314        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        // Set partitioned fields to prevent quadratic behavior.
341        // Ensure all row groups are partitioned.
342        row_groups
343            .into_par_iter()
344            .map(|(md, slice, row_count_start)| {
345                if slice.1 == 0 {
346                    return Ok(None);
347                }
348                // test we don't read the parquet file if this env var is set
349                #[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    // Fast path.
416    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    // @GB: I did quite some analysis on this.
483    //
484    // Pre-filtered and Post-filtered can both be faster in certain scenarios.
485    //
486    // - Pre-filtered is faster when there is some amount of clustering or
487    // sorting involved or if the number of values selected is small.
488    // - Post-filtering is faster when the predicate selects a somewhat random
489    // elements throughout the row group.
490    //
491    // The following is a heuristic value to try and estimate which one is
492    // faster. Essentially, it sees how many times it needs to switch between
493    // skipping items and collecting items and compares it against the number
494    // of values that it will collect.
495    //
496    // Closer to 0: pre-filtering is probably better.
497    // Closer to 1: post-filtering is probably better.
498    (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                // Prefiltering is only expensive for nested types so we make the cut-off quite
522                // high.
523                let is_nested = dtype.is_nested();
524
525                // We empirically selected these numbers.
526                !is_nested && prefilter_cost <= 0.01
527            },
528            Self::Pre => true,
529            Self::Post => false,
530        }
531    }
532}