Skip to main content

common_datasource/file_format/
parquet.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15pub 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
88/// Returns a AsyncFileReader factory
89impl 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    /// Must initialize the reader, or throw an error from the future.
134    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            // Safety: Must initialized
162            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
207/// Output the stream to a parquet file.
208///
209/// Returns number of rows written.
210pub 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}