1use std::io::{Read, Seek};
36use std::path::PathBuf;
37
38use arrow::datatypes::{ArrowSchemaRef, Metadata};
39use arrow::io::ipc::read::{self, get_row_count};
40use arrow::record_batch::RecordBatch;
41use polars_core::prelude::*;
42use polars_utils::bool::UnsafeBool;
43use polars_utils::pl_str::PlRefStr;
44#[cfg(feature = "serde")]
45use serde::{Deserialize, Serialize};
46
47use crate::RowIndex;
48use crate::hive::materialize_hive_partitions;
49use crate::mmap::MmapBytesReader;
50use crate::predicates::PhysicalIoExpr;
51use crate::prelude::*;
52use crate::shared::{ArrowReader, finish_reader};
53
54#[derive(Clone, Debug, PartialEq, Hash)]
55#[cfg_attr(feature = "serde", derive(Serialize, Deserialize))]
56#[cfg_attr(feature = "dsl-schema", derive(schemars::JsonSchema))]
57pub struct IpcScanOptions {
58 #[cfg_attr(feature = "serde", serde(default))]
60 pub record_batch_statistics: bool,
61 #[cfg_attr(feature = "serde", serde(default))]
62 pub checked: UnsafeBool,
63}
64
65#[expect(clippy::derivable_impls)]
66impl Default for IpcScanOptions {
67 fn default() -> Self {
68 Self {
69 record_batch_statistics: false,
70 checked: Default::default(),
71 }
72 }
73}
74
75#[must_use]
92pub struct IpcReader<R: MmapBytesReader> {
93 pub(super) reader: R,
95 rechunk: bool,
97 pub(super) n_rows: Option<usize>,
98 pub(super) projection: Option<Vec<usize>>,
99 pub(crate) columns: Option<Vec<String>>,
100 hive_partition_columns: Option<Vec<Series>>,
101 include_file_path: Option<(PlSmallStr, PlRefStr)>,
102 pub(super) row_index: Option<RowIndex>,
103 memory_map: memory_mapped_hidden::Key,
105 metadata: Option<read::FileMetadata>,
106 schema: Option<ArrowSchemaRef>,
107}
108
109mod memory_mapped_hidden {
110 use std::path::PathBuf;
111
112 #[derive(Default)]
113 pub struct Key {
114 inner: Option<PathBuf>,
115 }
116
117 impl Key {
118 pub unsafe fn new(inner: Option<PathBuf>) -> Self {
122 Self { inner }
123 }
124
125 pub fn is_set(&self) -> bool {
126 self.inner.is_some()
127 }
128 }
129}
130
131fn check_mmap_err(err: PolarsError) -> PolarsResult<()> {
132 if let PolarsError::ComputeError(s) = &err {
133 if s.as_ref() == "memory_map can only be done on uncompressed IPC files" {
134 eprintln!(
135 "Could not memory_map compressed IPC file, defaulting to normal read. \
136 Toggle off 'memory_map' to silence this warning."
137 );
138 return Ok(());
139 }
140 }
141 Err(err)
142}
143
144impl<R: MmapBytesReader> IpcReader<R> {
145 fn get_metadata(&mut self) -> PolarsResult<&read::FileMetadata> {
146 if self.metadata.is_none() {
147 let metadata = read::read_file_metadata(&mut self.reader)?;
148 self.schema = Some(metadata.schema.clone());
149 self.metadata = Some(metadata);
150 }
151 Ok(self.metadata.as_ref().unwrap())
152 }
153
154 pub fn schema(&mut self) -> PolarsResult<ArrowSchemaRef> {
156 self.get_metadata()?;
157 Ok(self.schema.as_ref().unwrap().clone())
158 }
159
160 pub fn custom_metadata(&mut self) -> PolarsResult<Option<Arc<Metadata>>> {
162 self.get_metadata()?;
163 Ok(self
164 .metadata
165 .as_ref()
166 .and_then(|meta| meta.custom_schema_metadata.clone()))
167 }
168
169 pub fn with_n_rows(mut self, num_rows: Option<usize>) -> Self {
171 self.n_rows = num_rows;
172 self
173 }
174
175 pub fn with_columns(mut self, columns: Option<Vec<String>>) -> Self {
177 self.columns = columns;
178 self
179 }
180
181 pub fn with_hive_partition_columns(mut self, columns: Option<Vec<Series>>) -> Self {
182 self.hive_partition_columns = columns;
183 self
184 }
185
186 pub fn with_include_file_path(
187 mut self,
188 include_file_path: Option<(PlSmallStr, PlRefStr)>,
189 ) -> Self {
190 self.include_file_path = include_file_path;
191 self
192 }
193
194 pub fn with_row_index(mut self, row_index: Option<RowIndex>) -> Self {
196 self.row_index = row_index;
197 self
198 }
199
200 pub fn with_projection(mut self, projection: Option<Vec<usize>>) -> Self {
203 self.projection = projection;
204 self
205 }
206
207 pub unsafe fn memory_mapped(mut self, path_buf: Option<PathBuf>) -> Self {
214 self.memory_map = unsafe { memory_mapped_hidden::Key::new(path_buf) };
215 self
216 }
217
218 #[cfg(feature = "lazy")]
220 pub fn finish_with_scan_ops(
221 mut self,
222 predicate: Option<Arc<dyn PhysicalIoExpr>>,
223 verbose: bool,
224 ) -> PolarsResult<DataFrame> {
225 if self.memory_map.is_set() && self.reader.to_file().is_some() {
226 if verbose {
227 eprintln!("memory map ipc file")
228 }
229 unsafe {
232 match self.finish_memmapped(predicate.clone()) {
233 Ok(df) => return Ok(df),
234 Err(err) => check_mmap_err(err)?,
235 }
236 }
237 }
238 let rechunk = self.rechunk;
239 let metadata = read::read_file_metadata(&mut self.reader)?;
240
241 if let Some(columns) = &self.columns {
245 self.projection = Some(columns_to_projection(columns, &metadata.schema)?);
246 }
247
248 let schema = if let Some(projection) = &self.projection {
249 Arc::new(apply_projection(&metadata.schema, projection))
250 } else {
251 metadata.schema.clone()
252 };
253
254 let reader = read::FileReader::new(self.reader, metadata, self.projection, self.n_rows);
255
256 finish_reader(reader, rechunk, None, predicate, &schema, self.row_index)
257 }
258}
259
260impl<R: MmapBytesReader> ArrowReader for read::FileReader<R>
261where
262 R: Read + Seek,
263{
264 fn next_record_batch(&mut self) -> PolarsResult<Option<RecordBatch>> {
265 self.next().map_or(Ok(None), |v| v.map(Some))
266 }
267}
268
269impl<R: MmapBytesReader> SerReader<R> for IpcReader<R> {
270 fn new(reader: R) -> Self {
271 IpcReader {
272 reader,
273 rechunk: true,
274 n_rows: None,
275 columns: None,
276 hive_partition_columns: None,
277 include_file_path: None,
278 projection: None,
279 row_index: None,
280 memory_map: Default::default(),
281 metadata: None,
282 schema: None,
283 }
284 }
285
286 fn set_rechunk(mut self, rechunk: bool) -> Self {
287 self.rechunk = rechunk;
288 self
289 }
290
291 fn finish(mut self) -> PolarsResult<DataFrame> {
292 let reader_schema = if let Some(ref schema) = self.schema {
293 schema.clone()
294 } else {
295 self.get_metadata()?.schema.clone()
296 };
297 let reader_schema = reader_schema.as_ref();
298
299 let hive_partition_columns = self.hive_partition_columns.take();
300 let include_file_path = self.include_file_path.take();
301
302 let mut df = (|| {
305 if self.projection.as_ref().is_some_and(|x| x.is_empty()) {
306 let row_count = if let Some(v) = self.n_rows {
307 v
308 } else {
309 get_row_count(&mut self.reader)? as usize
310 };
311 let mut df = DataFrame::empty_with_height(row_count);
312
313 if let Some(ri) = &self.row_index {
314 unsafe { df.with_row_index_mut(ri.name.clone(), Some(ri.offset)) };
315 }
316 return PolarsResult::Ok(df);
317 }
318
319 if self.memory_map.is_set() && self.reader.to_file().is_some() {
320 unsafe {
323 match self.finish_memmapped(None) {
324 Ok(df) => {
325 return Ok(df);
326 },
327 Err(err) => check_mmap_err(err)?,
328 }
329 }
330 }
331 let rechunk = self.rechunk;
332 let schema = self.get_metadata()?.schema.clone();
333
334 if let Some(columns) = &self.columns {
335 let prj = columns_to_projection(columns, schema.as_ref())?;
336 self.projection = Some(prj);
337 }
338
339 let schema = if let Some(projection) = &self.projection {
340 Arc::new(apply_projection(schema.as_ref(), projection))
341 } else {
342 schema
343 };
344
345 let metadata = self.get_metadata()?.clone();
346
347 let ipc_reader =
348 read::FileReader::new(self.reader, metadata, self.projection, self.n_rows);
349 let df = finish_reader(ipc_reader, rechunk, None, None, &schema, self.row_index)?;
350 Ok(df)
351 })()?;
352
353 if let Some(hive_cols) = hive_partition_columns {
354 materialize_hive_partitions(&mut df, reader_schema, Some(hive_cols.as_slice()));
355 };
356
357 if let Some((col, value)) = include_file_path {
358 unsafe {
359 df.push_column_unchecked(Column::new_scalar(
360 col,
361 Scalar::new(
362 DataType::String,
363 AnyValue::StringOwned(value.as_str().into()),
364 ),
365 df.height(),
366 ))
367 };
368 }
369
370 Ok(df)
371 }
372}