1use 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
33const DEFAULT_PREFETCH_SIZE: u64 = 64 * 1024;
35
36pub struct MetadataLoader<'a> {
37 object_store: ObjectStore,
39 file_path: &'a str,
41 file_size: u64,
43 page_index_policy: PageIndexPolicy,
44}
45
46impl<'a> MetadataLoader<'a> {
47 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 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 .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
113fn 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 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, ®ion_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, ®ion_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, ®ion_metadata));
350 }
351}