Skip to main content

polars_io/parquet/read/
async_impl.rs

1//! Read parquet files in parallel from the Object Store without a third party crate.
2
3use object_store::path::Path as ObjectPath;
4use polars_arrow::datatypes::ArrowSchemaRef;
5use polars_buffer::Buffer;
6use polars_core::prelude::*;
7use polars_parquet::parquet::error::ParquetError;
8use polars_parquet::parquet::read::{deserialize_metadata, deserialize_num_rows};
9use polars_parquet::parquet::{FOOTER_SIZE, PARQUET_MAGIC};
10use polars_utils::pl_path::PlRefPath;
11
12use crate::cloud::concurrency_config::FetchConfig;
13use crate::cloud::{
14    CloudLocation, CloudOptions, PolarsObjectStore, build_object_store, object_path_from_str,
15};
16use crate::configs::cloud_footer_read_size;
17use crate::parquet::metadata::FileMetadataRef;
18
19pub struct ParquetObjectStore {
20    store: PolarsObjectStore,
21    path: ObjectPath,
22    metadata: Option<FileMetadataRef>,
23    schema: Option<ArrowSchemaRef>,
24}
25
26impl ParquetObjectStore {
27    pub async fn from_uri(
28        uri: PlRefPath,
29        options: Option<&CloudOptions>,
30        metadata: Option<FileMetadataRef>,
31    ) -> PolarsResult<Self> {
32        let (CloudLocation { prefix, .. }, store) = build_object_store(uri, options, false).await?;
33        let path = object_path_from_str(&prefix)?;
34
35        Ok(ParquetObjectStore {
36            store,
37            path,
38            metadata,
39            schema: None,
40        })
41    }
42
43    /// Number of rows in the parquet file.
44    pub async fn num_rows(&mut self) -> PolarsResult<usize> {
45        let metadata = self.get_metadata().await?;
46        Ok(metadata.num_rows)
47    }
48
49    /// Fetch and memoize the metadata of the parquet file.
50    pub async fn get_metadata(&mut self) -> PolarsResult<&FileMetadataRef> {
51        if self.metadata.is_none() {
52            let footer = fetch_footer_bytes(&self.store, &self.path).await?;
53            self.metadata = Some(Arc::new(deserialize_metadata(footer)?));
54        }
55        Ok(self.metadata.as_ref().unwrap())
56    }
57
58    /// Decode only `FileMetaData.num_rows` from the remote footer.
59    /// Not memoized. Used by `RowCounts` resolve mode.
60    pub async fn num_rows_only(&mut self) -> PolarsResult<i64> {
61        let footer = fetch_footer_bytes(&self.store, &self.path).await?;
62        Ok(deserialize_num_rows(footer)?)
63    }
64
65    pub async fn schema(&mut self) -> PolarsResult<ArrowSchemaRef> {
66        self.schema = Some(match self.schema.as_ref() {
67            Some(schema) => Arc::clone(schema),
68            None => {
69                let metadata = self.get_metadata().await?;
70                let arrow_schema = polars_parquet::arrow::read::infer_schema(metadata)?;
71                Arc::new(arrow_schema)
72            },
73        });
74
75        Ok(self.schema.clone().unwrap())
76    }
77}
78
79fn read_n<const N: usize>(reader: &mut &[u8]) -> Option<[u8; N]> {
80    if N <= reader.len() {
81        let (head, tail) = reader.split_at(N);
82        *reader = tail;
83        Some(head.try_into().unwrap())
84    } else {
85        None
86    }
87}
88
89fn read_i32le(reader: &mut &[u8]) -> Option<i32> {
90    read_n(reader).map(i32::from_le_bytes)
91}
92
93/// Speculatively read `cloud_footer_read_size()` bytes from the tail. If the
94/// footer fits in the prefetch (the common case), we're done in one range
95/// request; otherwise re-fetch the full footer. Mirrors the sync
96/// `fetch_footer_buf` strategy.
97async fn fetch_footer_bytes(
98    store: &PolarsObjectStore,
99    path: &ObjectPath,
100) -> PolarsResult<Buffer<u8>> {
101    let out_of_spec = |msg: &str| ParquetError::OutOfSpec(msg.to_string());
102
103    let (prefetched, file_byte_length) = store
104        .get_suffix(path, cloud_footer_read_size(), FetchConfig::random_access())
105        .await?;
106
107    if prefetched.len() < FOOTER_SIZE as usize {
108        return Err(out_of_spec("not enough bytes to contain parquet footer").into());
109    }
110
111    // Trailing 8 bytes: footer size (i32 LE) + magic.
112    let footer_byte_length: usize = {
113        let tail_start = prefetched.len() - FOOTER_SIZE as usize;
114        let reader = &mut &prefetched.as_ref()[tail_start..];
115        let footer_byte_size = read_i32le(reader).unwrap();
116        let magic = read_n(reader).unwrap();
117        debug_assert!(reader.is_empty());
118        if magic != PARQUET_MAGIC {
119            return Err(out_of_spec("incorrect magic in parquet footer").into());
120        }
121        footer_byte_size
122            .try_into()
123            .map_err(|_| out_of_spec("negative footer byte length"))?
124    };
125
126    let footer_len = FOOTER_SIZE as usize + footer_byte_length;
127    if footer_len <= prefetched.len() {
128        // Common case: footer already in the prefetch; zero extra round trips.
129        let footer = prefetched.clone().sliced((prefetched.len() - footer_len)..);
130
131        // The footer is held for the lifetime of the metadata, so copy it out rather than
132        // pin the whole prefetch for a fraction of its bytes.
133        return Ok(if prefetched.len() >= 2 * footer_len {
134            Buffer::from_vec(footer.to_vec())
135        } else {
136            footer
137        });
138    }
139
140    // Fallback: footer larger than the prefetch; re-fetch the full footer.
141    store
142        .get_range(
143            path,
144            file_byte_length
145                .checked_sub(footer_len)
146                .ok_or_else(|| out_of_spec("not enough bytes to contain parquet footer"))?
147                ..file_byte_length,
148            FetchConfig::random_access(),
149        )
150        .await
151}