polars_lazy/scan/
parquet.rs1use polars_buffer::Buffer;
2use polars_core::prelude::*;
3use polars_io::cloud::CloudOptions;
4use polars_io::parquet::read::ParallelStrategy;
5use polars_io::prelude::ParquetOptions;
6use polars_io::{HiveOptions, RowIndex};
7use polars_utils::pl_path::PlRefPath;
8use polars_utils::slice_enum::Slice;
9
10use crate::prelude::*;
11
12#[derive(Clone)]
13pub struct ScanArgsParquet {
14 pub n_rows: Option<usize>,
15 pub parallel: ParallelStrategy,
16 pub row_index: Option<RowIndex>,
17 pub cloud_options: Option<CloudOptions>,
18 pub hive_options: HiveOptions,
19 pub use_statistics: bool,
20 pub schema: Option<SchemaRef>,
21 pub low_memory: bool,
22 pub rechunk: bool,
23 pub cache: bool,
24 pub glob: bool,
26 pub include_file_paths: Option<PlSmallStr>,
27 pub allow_missing_columns: bool,
28}
29
30impl Default for ScanArgsParquet {
31 fn default() -> Self {
32 Self {
33 n_rows: None,
34 parallel: Default::default(),
35 row_index: None,
36 cloud_options: None,
37 hive_options: Default::default(),
38 use_statistics: true,
39 schema: None,
40 rechunk: false,
41 low_memory: false,
42 cache: true,
43 glob: true,
44 include_file_paths: None,
45 allow_missing_columns: false,
46 }
47 }
48}
49
50#[derive(Clone)]
51struct LazyParquetReader {
52 args: ScanArgsParquet,
53 sources: ScanSources,
54}
55
56impl LazyParquetReader {
57 fn new(args: ScanArgsParquet) -> Self {
58 Self {
59 args,
60 sources: ScanSources::default(),
61 }
62 }
63}
64
65impl LazyFileListReader for LazyParquetReader {
66 fn finish(self) -> PolarsResult<LazyFrame> {
68 let parquet_options = ParquetOptions {
69 schema: self.args.schema,
70 parallel: self.args.parallel,
71 low_memory: self.args.low_memory,
72 use_statistics: self.args.use_statistics,
73 };
74
75 let unified_scan_args = UnifiedScanArgs {
76 schema: None,
77 cloud_options: self.args.cloud_options,
78 hive_options: self.args.hive_options,
79 rechunk: self.args.rechunk,
80 cache: self.args.cache,
81 glob: self.args.glob,
82 expand_paths: true,
83 hidden_file_prefix: None,
84 projection: None,
85 column_mapping: None,
86 default_values: None,
87 row_index: None,
89 pre_slice: self
90 .args
91 .n_rows
92 .map(|len| Slice::Positive { offset: 0, len }),
93 cast_columns_policy: CastColumnsPolicy::ERROR_ON_MISMATCH,
94 missing_columns_policy: if self.args.allow_missing_columns {
95 MissingColumnsPolicy::Insert
96 } else {
97 MissingColumnsPolicy::Raise
98 },
99 extra_columns_policy: ExtraColumnsPolicy::Raise,
100 include_file_paths: self.args.include_file_paths,
101 deletion_files: None,
102 table_statistics: None,
103 row_count: None,
104 source_sizes: None,
105 resolve_heavy_sources: None,
106 };
107
108 let mut lf: LazyFrame =
109 DslBuilder::scan_parquet(self.sources, parquet_options, unified_scan_args)?
110 .build()
111 .into();
112
113 if let Some(row_index) = self.args.row_index {
115 lf = lf.with_row_index(row_index.name, Some(row_index.offset))
116 }
117
118 Ok(lf)
119 }
120
121 fn glob(&self) -> bool {
122 self.args.glob
123 }
124
125 fn finish_no_glob(self) -> PolarsResult<LazyFrame> {
126 unreachable!();
127 }
128
129 fn sources(&self) -> &ScanSources {
130 &self.sources
131 }
132
133 fn with_sources(mut self, sources: ScanSources) -> Self {
134 self.sources = sources;
135 self
136 }
137
138 fn with_n_rows(mut self, n_rows: impl Into<Option<usize>>) -> Self {
139 self.args.n_rows = n_rows.into();
140 self
141 }
142
143 fn with_row_index(mut self, row_index: impl Into<Option<RowIndex>>) -> Self {
144 self.args.row_index = row_index.into();
145 self
146 }
147
148 fn rechunk(&self) -> bool {
149 self.args.rechunk
150 }
151
152 fn with_rechunk(mut self, toggle: bool) -> Self {
153 self.args.rechunk = toggle;
154 self
155 }
156
157 fn cloud_options(&self) -> Option<&CloudOptions> {
158 self.args.cloud_options.as_ref()
159 }
160
161 fn n_rows(&self) -> Option<usize> {
162 self.args.n_rows
163 }
164
165 fn row_index(&self) -> Option<&RowIndex> {
166 self.args.row_index.as_ref()
167 }
168}
169
170impl LazyFrame {
171 pub fn scan_parquet(path: PlRefPath, args: ScanArgsParquet) -> PolarsResult<Self> {
173 Self::scan_parquet_sources(ScanSources::Paths(Buffer::from_iter([path])), args)
174 }
175
176 pub fn scan_parquet_sources(sources: ScanSources, args: ScanArgsParquet) -> PolarsResult<Self> {
178 LazyParquetReader::new(args).with_sources(sources).finish()
179 }
180
181 pub fn scan_parquet_files(
183 paths: Buffer<PlRefPath>,
184 args: ScanArgsParquet,
185 ) -> PolarsResult<Self> {
186 Self::scan_parquet_sources(ScanSources::Paths(paths), args)
187 }
188}