polars_io/parquet/read/
async_impl.rs1use 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 pub async fn num_rows(&mut self) -> PolarsResult<usize> {
45 let metadata = self.get_metadata().await?;
46 Ok(metadata.num_rows)
47 }
48
49 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 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
93async 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 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 let footer = prefetched.clone().sliced((prefetched.len() - footer_len)..);
130
131 return Ok(if prefetched.len() >= 2 * footer_len {
134 Buffer::from_vec(footer.to_vec())
135 } else {
136 footer
137 });
138 }
139
140 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}