Skip to main content

polars_io/ipc/
ipc_file.rs

1//! # (De)serializing Arrows IPC format.
2//!
3//! Arrow IPC is a [binary format](https://arrow.apache.org/docs/python/ipc.html).
4//! It is the recommended way to serialize and deserialize Polars DataFrames as this is most true
5//! to the data schema.
6//!
7//! ## Example
8//!
9//! ```rust
10//! use polars_core::prelude::*;
11//! use polars_io::prelude::*;
12//! use std::io::Cursor;
13//!
14//!
15//! let s0 = Column::new("days".into(), &[0, 1, 2, 3, 4]);
16//! let s1 = Column::new("temp".into(), &[22.1, 19.9, 7., 2., 3.]);
17//! let mut df = DataFrame::new_infer_height(vec![s0, s1]).unwrap();
18//!
19//! // Create an in memory file handler.
20//! // Vec<u8>: Read + Write
21//! // Cursor<T>: Seek
22//!
23//! let mut buf: Cursor<Vec<u8>> = Cursor::new(Vec::new());
24//!
25//! // write to the in memory buffer
26//! IpcWriter::new(&mut buf).finish(&mut df).expect("ipc writer");
27//!
28//! // reset the buffers index after writing to the beginning of the buffer
29//! buf.set_position(0);
30//!
31//! // read the buffer into a DataFrame
32//! let df_read = IpcReader::new(buf).finish().unwrap();
33//! assert!(df.equals(&df_read));
34//! ```
35use 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    /// Read StatisticsFlags from the record batch custom metadata.
59    #[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/// Read Arrows IPC format into a DataFrame
76///
77/// # Example
78/// ```
79/// use polars_core::prelude::*;
80/// use std::fs::File;
81/// use polars_io::ipc::IpcReader;
82/// use polars_io::SerReader;
83///
84/// fn example() -> PolarsResult<DataFrame> {
85///     let file = File::open("file.ipc").expect("file not found");
86///
87///     IpcReader::new(file)
88///         .finish()
89/// }
90/// ```
91#[must_use]
92pub struct IpcReader<R: MmapBytesReader> {
93    /// File or Stream object
94    pub(super) reader: R,
95    /// Aggregates chunks afterwards to a single chunk.
96    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    // Stores the as key semaphore to make sure we don't write to the memory mapped file.
104    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        /// # Safety
119        /// The users guarantees that the arrow data in the given path is valid
120        /// and remains valid throughout memory mapping.
121        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    /// Get arrow schema of the Ipc File.
155    pub fn schema(&mut self) -> PolarsResult<ArrowSchemaRef> {
156        self.get_metadata()?;
157        Ok(self.schema.as_ref().unwrap().clone())
158    }
159
160    /// Get schema-level custom metadata of the Ipc file
161    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    /// Stop reading when `n` rows are read.
170    pub fn with_n_rows(mut self, num_rows: Option<usize>) -> Self {
171        self.n_rows = num_rows;
172        self
173    }
174
175    /// Columns to select/ project
176    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    /// Add a row index column.
195    pub fn with_row_index(mut self, row_index: Option<RowIndex>) -> Self {
196        self.row_index = row_index;
197        self
198    }
199
200    /// Set the reader's column projection. This counts from 0, meaning that
201    /// `vec![0, 4]` would select the 1st and 5th column.
202    pub fn with_projection(mut self, projection: Option<Vec<usize>>) -> Self {
203        self.projection = projection;
204        self
205    }
206
207    /// Set if the file is to be memory_mapped. Only works with uncompressed files.
208    /// The file name must be passed to register the memory mapped file.
209    ///
210    /// # Safety
211    /// The users guarantees that the arrow data in the given path is valid
212    /// and remains valid throughout memory mapping.
213    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    // todo! hoist to lazy crate
219    #[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            // # Safety
230            // Can only be set by a user that guarantees correct arrow data.
231            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        // NOTE: For some code paths this already happened. See
242        // https://github.com/pola-rs/polars/pull/14984#discussion_r1520125000
243        // where this was introduced.
244        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        // In case only hive columns are projected, the df would be empty, but we need the row count
303        // of the file in order to project the correct number of rows for the hive columns.
304        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                // Safety:
321                // Can only be set by user that guarantees valid arrow data.
322                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}