Skip to main content

mito2/sst/parquet/
metadata.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
15use std::result::Result as StdResult;
16use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
17
18use bytes::Bytes;
19use futures::FutureExt;
20use futures::future::BoxFuture;
21use object_store::ObjectStore;
22use parquet::arrow::async_reader::MetadataFetch;
23use parquet::errors::{ParquetError, Result as ParquetResult};
24use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader};
25use parquet::file::statistics::Statistics;
26use snafu::{IntoError as _, ResultExt};
27use store_api::metadata::RegionMetadata;
28use store_api::storage::consts::PRIMARY_KEY_COLUMN_NAME;
29
30use crate::error::{self, Result};
31use crate::sst::parquet::reader::MetadataCacheMetrics;
32
33/// The estimated size of the footer and metadata need to read from the end of parquet file.
34const DEFAULT_PREFETCH_SIZE: u64 = 64 * 1024;
35
36pub struct MetadataLoader<'a> {
37    // An object store that supports async read
38    object_store: ObjectStore,
39    // The path of parquet file
40    file_path: &'a str,
41    // The size of parquet file
42    file_size: u64,
43    page_index_policy: PageIndexPolicy,
44}
45
46impl<'a> MetadataLoader<'a> {
47    /// Create a new parquet metadata loader.
48    pub fn new(
49        object_store: ObjectStore,
50        file_path: &'a str,
51        file_size: u64,
52    ) -> MetadataLoader<'a> {
53        Self {
54            object_store,
55            file_path,
56            file_size,
57            page_index_policy: Default::default(),
58        }
59    }
60
61    pub(crate) fn with_page_index_policy(&mut self, page_index_policy: PageIndexPolicy) {
62        self.page_index_policy = page_index_policy;
63    }
64
65    /// Get the size of parquet file. If file_size is 0, stat the object store to get the size.
66    async fn get_file_size(&self) -> Result<u64> {
67        let file_size = match self.file_size {
68            0 => self
69                .object_store
70                .stat(self.file_path)
71                .await
72                .context(error::OpenDalSnafu)?
73                .content_length(),
74            other => other,
75        };
76        Ok(file_size)
77    }
78
79    pub async fn load(&self, cache_metrics: &mut MetadataCacheMetrics) -> Result<ParquetMetaData> {
80        let path = self.file_path;
81        let file_size = self.get_file_size().await?;
82        let reader = ParquetMetaDataReader::new()
83            .with_prefetch_hint(Some(DEFAULT_PREFETCH_SIZE as usize))
84            // Mito uses offset indexes to translate row selections into byte ranges. It does not
85            // consume Parquet column indexes, so decoding them only bloats the metadata cache.
86            .with_column_index_policy(PageIndexPolicy::Skip)
87            .with_offset_index_policy(self.page_index_policy);
88
89        let num_reads = AtomicUsize::new(0);
90        let bytes_read = AtomicU64::new(0);
91        let fetch = ObjectStoreFetch {
92            object_store: &self.object_store,
93            file_path: self.file_path,
94            num_reads: &num_reads,
95            bytes_read: &bytes_read,
96        };
97
98        let metadata = reader
99            .load_and_finish(fetch, file_size)
100            .await
101            .map_err(|e| match unbox_external_error(e) {
102                Ok(os_err) => error::OpenDalSnafu {}.into_error(os_err),
103                Err(parquet_err) => error::ReadParquetSnafu { path }.into_error(parquet_err),
104            })?;
105
106        cache_metrics.num_reads = num_reads.into_inner();
107        cache_metrics.bytes_read = bytes_read.into_inner();
108
109        Ok(metadata)
110    }
111}
112
113/// Unpack ParquetError to get object_store::Error if possible.
114fn unbox_external_error(e: ParquetError) -> StdResult<object_store::Error, ParquetError> {
115    match e {
116        ParquetError::External(boxed_err) => match boxed_err.downcast::<object_store::Error>() {
117            Ok(os_err) => Ok(*os_err),
118            Err(parquet_error) => Err(ParquetError::External(parquet_error)),
119        },
120        other => Err(other),
121    }
122}
123
124pub(crate) fn extract_primary_key_range(
125    parquet_meta: &ParquetMetaData,
126    region_metadata: &RegionMetadata,
127) -> Option<(Bytes, Bytes)> {
128    if region_metadata.primary_key.is_empty() {
129        return None;
130    }
131
132    let pk_column_idx = parquet_meta
133        .file_metadata()
134        .schema_descr()
135        .columns()
136        .iter()
137        .position(|column| column.name() == PRIMARY_KEY_COLUMN_NAME)?;
138
139    let mut min: Option<Bytes> = None;
140    let mut max: Option<Bytes> = None;
141
142    for row_group in parquet_meta.row_groups() {
143        let Statistics::ByteArray(stats) = row_group.column(pk_column_idx).statistics()? else {
144            return None;
145        };
146
147        // Truncated bounds can end at a field boundary and masquerade as an
148        // older Dense schema. Only complete endpoints can be schema-normalized.
149        if !stats.min_is_exact() || !stats.max_is_exact() {
150            return None;
151        }
152        let row_group_min = Bytes::copy_from_slice(stats.min_bytes_opt()?);
153        let row_group_max = Bytes::copy_from_slice(stats.max_bytes_opt()?);
154        min = Some(match min {
155            Some(current) => current.min(row_group_min),
156            None => row_group_min,
157        });
158        max = Some(match max {
159            Some(current) => current.max(row_group_max),
160            None => row_group_max,
161        });
162    }
163
164    min.zip(max)
165}
166
167struct ObjectStoreFetch<'a> {
168    object_store: &'a ObjectStore,
169    file_path: &'a str,
170    num_reads: &'a AtomicUsize,
171    bytes_read: &'a AtomicU64,
172}
173
174impl MetadataFetch for ObjectStoreFetch<'_> {
175    fn fetch(&mut self, range: std::ops::Range<u64>) -> BoxFuture<'_, ParquetResult<Bytes>> {
176        let bytes_to_read = range.end - range.start;
177        async move {
178            let data = self
179                .object_store
180                .read_with(self.file_path)
181                .range(range)
182                .await
183                .map_err(|e| ParquetError::External(Box::new(e)))?;
184            self.num_reads.fetch_add(1, Ordering::Relaxed);
185            self.bytes_read.fetch_add(bytes_to_read, Ordering::Relaxed);
186            Ok(data.to_bytes())
187        }
188        .boxed()
189    }
190}
191
192#[cfg(test)]
193mod tests {
194    use std::sync::Arc;
195
196    use datatypes::arrow::array::{
197        ArrayRef, BinaryArray, DictionaryArray, Int64Array, UInt32Array,
198    };
199    use datatypes::arrow::datatypes::{DataType as ArrowDataType, Field, Schema};
200    use datatypes::arrow::record_batch::RecordBatch;
201    use object_store::services::Memory;
202    use parquet::arrow::ArrowWriter;
203    use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
204    use parquet::file::metadata::{KeyValue, ParquetMetaData};
205    use parquet::file::properties::{EnabledStatistics, WriterProperties};
206
207    use super::*;
208    use crate::sst::parquet::PARQUET_METADATA_KEY;
209    use crate::test_util::sst_util::sst_region_metadata;
210
211    fn build_test_parquet_bytes(
212        include_primary_key: bool,
213        primary_keys: &[&[u8]],
214        row_group_sizes: &[usize],
215        stats_enabled: EnabledStatistics,
216    ) -> Vec<u8> {
217        let total_rows = row_group_sizes.iter().sum::<usize>();
218        let mut fields = vec![Field::new("field", ArrowDataType::Int64, true)];
219        let mut columns: Vec<ArrayRef> =
220            vec![Arc::new(Int64Array::from_iter_values(0..total_rows as i64))];
221        if include_primary_key {
222            assert_eq!(total_rows, primary_keys.len());
223            fields.push(Field::new(
224                "__primary_key",
225                ArrowDataType::Dictionary(
226                    Box::new(ArrowDataType::UInt32),
227                    Box::new(ArrowDataType::Binary),
228                ),
229                false,
230            ));
231            let values = Arc::new(BinaryArray::from_iter_values(primary_keys.iter().copied()));
232            let keys = UInt32Array::from_iter_values(0..primary_keys.len() as u32);
233            columns.push(Arc::new(DictionaryArray::new(keys, values)));
234        }
235
236        let schema = Arc::new(Schema::new(fields));
237        let region_metadata = Arc::new(sst_region_metadata());
238        let key_value = KeyValue::new(
239            PARQUET_METADATA_KEY.to_string(),
240            region_metadata.to_json().unwrap(),
241        );
242        let props = WriterProperties::builder()
243            .set_key_value_metadata(Some(vec![key_value]))
244            .set_statistics_enabled(stats_enabled)
245            .build();
246
247        let mut parquet_bytes = Vec::new();
248        let mut writer =
249            ArrowWriter::try_new(&mut parquet_bytes, schema.clone(), Some(props)).unwrap();
250        let mut offset = 0;
251        for row_group_size in row_group_sizes {
252            let batch = RecordBatch::try_new(
253                schema.clone(),
254                columns
255                    .iter()
256                    .map(|column| column.slice(offset, *row_group_size))
257                    .collect(),
258            )
259            .unwrap();
260            writer.write(&batch).unwrap();
261            offset += row_group_size;
262        }
263        writer.close().unwrap();
264
265        parquet_bytes
266    }
267
268    fn build_test_metadata(
269        include_primary_key: bool,
270        primary_keys: &[&[u8]],
271        row_group_sizes: &[usize],
272        stats_enabled: EnabledStatistics,
273    ) -> ParquetMetaData {
274        ParquetRecordBatchReaderBuilder::try_new(Bytes::from(build_test_parquet_bytes(
275            include_primary_key,
276            primary_keys,
277            row_group_sizes,
278            stats_enabled,
279        )))
280        .unwrap()
281        .metadata()
282        .as_ref()
283        .clone()
284    }
285
286    #[tokio::test]
287    async fn test_metadata_loader_only_loads_offset_indexes() {
288        let parquet_bytes = build_test_parquet_bytes(false, &[], &[4], EnabledStatistics::Page);
289        let file_size = parquet_bytes.len() as u64;
290        let file_path = "test.parquet";
291        let object_store = ObjectStore::new(Memory::default()).unwrap();
292        object_store.write(file_path, parquet_bytes).await.unwrap();
293
294        let mut loader = MetadataLoader::new(object_store, file_path, file_size);
295        loader.with_page_index_policy(PageIndexPolicy::Required);
296        let metadata = loader
297            .load(&mut MetadataCacheMetrics::default())
298            .await
299            .unwrap();
300
301        assert!(metadata.column_index().is_none());
302        assert!(metadata.offset_index().is_some());
303    }
304
305    #[test]
306    fn test_extract_primary_key_range_returns_none_when_column_absent() {
307        let metadata = build_test_metadata(false, &[], &[1], EnabledStatistics::Page);
308        let region_metadata = sst_region_metadata();
309
310        assert_eq!(None, extract_primary_key_range(&metadata, &region_metadata));
311    }
312
313    #[test]
314    fn test_extract_primary_key_range_folds_row_group_stats() {
315        let metadata = build_test_metadata(
316            true,
317            &[b"bbb", b"ccc", b"aaa", b"zzz"],
318            &[2, 2],
319            EnabledStatistics::Page,
320        );
321        let region_metadata = sst_region_metadata();
322
323        assert_eq!(
324            Some((Bytes::from_static(b"aaa"), Bytes::from_static(b"zzz"))),
325            extract_primary_key_range(&metadata, &region_metadata)
326        );
327    }
328
329    #[test]
330    fn test_extract_primary_key_range_rejects_truncated_statistics() {
331        let key = vec![b'a'; 1024];
332        let metadata = build_test_metadata(true, &[&key], &[1], EnabledStatistics::Page);
333        assert_eq!(
334            None,
335            extract_primary_key_range(&metadata, &sst_region_metadata())
336        );
337    }
338
339    #[test]
340    fn test_extract_primary_key_range_returns_none_when_any_rg_stats_missing() {
341        let metadata = build_test_metadata(
342            true,
343            &[b"bbb", b"ccc", b"aaa", b"zzz"],
344            &[2, 2],
345            EnabledStatistics::None,
346        );
347        let region_metadata = sst_region_metadata();
348
349        assert_eq!(None, extract_primary_key_range(&metadata, &region_metadata));
350    }
351}