common_datasource/file_format/
parquet.rs1pub mod packed_reader;
16
17use std::result;
18use std::sync::Arc;
19
20use arrow::record_batch::RecordBatch;
21use arrow_schema::Schema;
22use async_trait::async_trait;
23use datafusion::datasource::physical_plan::ParquetFileReaderFactory;
24use datafusion::error::Result as DatafusionResult;
25use datafusion::parquet::arrow::async_reader::AsyncFileReader;
26use datafusion::parquet::arrow::{ArrowWriter, parquet_to_arrow_schema};
27use datafusion::parquet::errors::{ParquetError, Result as ParquetResult};
28use datafusion::parquet::file::metadata::{
29 PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader,
30};
31use datafusion::physical_plan::SendableRecordBatchStream;
32use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet;
33use datafusion_datasource::PartitionedFile;
34use futures::StreamExt;
35use futures::future::BoxFuture;
36use object_store::{FuturesAsyncReader, ObjectStore};
37use parquet::arrow::arrow_reader::ArrowReaderOptions;
38use snafu::ResultExt;
39use tokio_util::compat::{Compat, FuturesAsyncReadCompatExt};
40
41use crate::buffered_writer::{ArrowWriterCloser, DfRecordBatchEncoder};
42use crate::error::{self, Result};
43use crate::file_format::FileFormat;
44use crate::parquet_writer::ParquetFileWriter;
45use crate::share_buffer::SharedBuffer;
46
47#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
48pub struct ParquetFormat {}
49
50#[async_trait]
51impl FileFormat for ParquetFormat {
52 async fn infer_schema(&self, store: &ObjectStore, path: &str) -> Result<Schema> {
53 let meta = store
54 .stat(path)
55 .await
56 .context(error::ReadObjectSnafu { path })?;
57
58 let mut reader = store
59 .reader(path)
60 .await
61 .context(error::ReadObjectSnafu { path })?
62 .into_futures_async_read(0..meta.content_length())
63 .await
64 .context(error::ReadObjectSnafu { path })?
65 .compat();
66
67 let metadata = reader
68 .get_metadata(None)
69 .await
70 .context(error::ReadParquetSnafuSnafu)?;
71
72 let file_metadata = metadata.file_metadata();
73 let schema = parquet_to_arrow_schema(
74 file_metadata.schema_descr(),
75 file_metadata.key_value_metadata(),
76 )
77 .context(error::ParquetToSchemaSnafu)?;
78
79 Ok(schema)
80 }
81}
82
83#[derive(Debug, Clone)]
84pub struct DefaultParquetFileReaderFactory {
85 object_store: ObjectStore,
86}
87
88impl DefaultParquetFileReaderFactory {
90 pub fn new(object_store: ObjectStore) -> Self {
91 Self { object_store }
92 }
93}
94
95impl ParquetFileReaderFactory for DefaultParquetFileReaderFactory {
96 fn create_reader(
97 &self,
98 _partition_index: usize,
99 partitioned_file: PartitionedFile,
100 metadata_size_hint: Option<usize>,
101 _metrics: &ExecutionPlanMetricsSet,
102 ) -> DatafusionResult<Box<dyn AsyncFileReader + Send>> {
103 let path = partitioned_file.path().to_string();
104 let object_store = self.object_store.clone();
105
106 Ok(Box::new(LazyParquetFileReader::new(
107 object_store,
108 path,
109 metadata_size_hint,
110 )))
111 }
112}
113
114pub struct LazyParquetFileReader {
115 object_store: ObjectStore,
116 reader: Option<Compat<FuturesAsyncReader>>,
117 file_size: Option<u64>,
118 metadata_size_hint: Option<usize>,
119 path: String,
120}
121
122impl LazyParquetFileReader {
123 pub fn new(object_store: ObjectStore, path: String, metadata_size_hint: Option<usize>) -> Self {
124 LazyParquetFileReader {
125 object_store,
126 path,
127 reader: None,
128 file_size: None,
129 metadata_size_hint,
130 }
131 }
132
133 async fn maybe_initialize(&mut self) -> result::Result<(), object_store::Error> {
135 if self.reader.is_none() {
136 let meta = self.object_store.stat(&self.path).await?;
137 self.file_size = Some(meta.content_length());
138 let reader = self
139 .object_store
140 .reader(&self.path)
141 .await?
142 .into_futures_async_read(0..meta.content_length())
143 .await?
144 .compat();
145 self.reader = Some(reader);
146 }
147
148 Ok(())
149 }
150}
151
152impl AsyncFileReader for LazyParquetFileReader {
153 fn get_bytes(
154 &mut self,
155 range: std::ops::Range<u64>,
156 ) -> BoxFuture<'_, ParquetResult<bytes::Bytes>> {
157 Box::pin(async move {
158 self.maybe_initialize()
159 .await
160 .map_err(|e| ParquetError::External(Box::new(e)))?;
161 self.reader.as_mut().unwrap().get_bytes(range).await
163 })
164 }
165
166 fn get_metadata<'a>(
167 &'a mut self,
168 options: Option<&'a ArrowReaderOptions>,
169 ) -> BoxFuture<'a, parquet::errors::Result<Arc<ParquetMetaData>>> {
170 Box::pin(async move {
171 self.maybe_initialize()
172 .await
173 .map_err(|e| ParquetError::External(Box::new(e)))?;
174
175 let metadata_opts = options.map(|o| o.metadata_options().clone());
176 let column_index_policy =
177 options.map_or(PageIndexPolicy::Skip, |o| o.column_index_policy());
178 let offset_index_policy =
179 options.map_or(PageIndexPolicy::Skip, |o| o.offset_index_policy());
180 let metadata_reader = ParquetMetaDataReader::new()
181 .with_metadata_options(metadata_opts)
182 .with_column_index_policy(column_index_policy)
183 .with_offset_index_policy(offset_index_policy)
184 .with_prefetch_hint(self.metadata_size_hint);
185
186 let metadata = metadata_reader
187 .load_and_finish(self.reader.as_mut().unwrap(), self.file_size.unwrap())
188 .await?;
189 Ok(Arc::new(metadata))
190 })
191 }
192}
193
194impl DfRecordBatchEncoder for ArrowWriter<SharedBuffer> {
195 fn write(&mut self, batch: &RecordBatch) -> Result<()> {
196 self.write(batch).context(error::EncodeRecordBatchSnafu)
197 }
198}
199
200#[async_trait]
201impl ArrowWriterCloser for ArrowWriter<SharedBuffer> {
202 async fn close(self) -> Result<ParquetMetaData> {
203 self.close().context(error::EncodeRecordBatchSnafu)
204 }
205}
206
207pub async fn stream_to_parquet(
211 mut stream: SendableRecordBatchStream,
212 store: ObjectStore,
213 path: &str,
214 concurrency: usize,
215) -> Result<usize> {
216 let mut writer =
217 ParquetFileWriter::open(stream.schema(), store, path, concurrency, None).await?;
218 let mut rows_written = 0;
219 while let Some(batch) = stream.next().await {
220 let batch = batch.context(error::ReadRecordBatchSnafu)?;
221 let rows = batch.num_rows();
222 writer.write(batch, None).await?;
223 rows_written += rows;
224 }
225 writer.finish(None).await?;
226 Ok(rows_written)
227}
228
229#[cfg(test)]
230mod tests {
231 use common_test_util::find_workspace_path;
232
233 use super::*;
234 use crate::test_util::{format_schema, test_store};
235
236 fn test_data_root() -> String {
237 find_workspace_path("/src/common/datasource/tests/parquet")
238 .display()
239 .to_string()
240 }
241
242 #[tokio::test]
243 async fn infer_schema_basic() {
244 let json = ParquetFormat::default();
245 let store = test_store(&test_data_root());
246 let schema = json.infer_schema(&store, "basic.parquet").await.unwrap();
247 let formatted: Vec<_> = format_schema(schema);
248
249 assert_eq!(vec!["num: Int64: NULL", "str: Utf8: NULL"], formatted);
250 }
251}