Skip to main content

micromegas_analytics/lakehouse/
reader_factory.rs

1use super::metadata_cache::MetadataCache;
2use super::partition_metadata::load_partition_metadata;
3use bytes::Bytes;
4use datafusion::{
5    datasource::{
6        listing::PartitionedFile,
7        physical_plan::{ParquetFileMetrics, ParquetFileReaderFactory},
8    },
9    parquet::{
10        arrow::{arrow_reader::ArrowReaderOptions, async_reader::AsyncFileReader},
11        file::metadata::ParquetMetaData,
12    },
13    physical_plan::metrics::{Count, ExecutionPlanMetricsSet},
14};
15use futures::future::BoxFuture;
16use micromegas_tracing::prelude::*;
17use object_store::{ObjectStore, ObjectStoreExt, path::Path};
18use std::ops::Range;
19use std::sync::Arc;
20
21/// A custom [`ParquetFileReaderFactory`] that handles opening parquet files
22/// from object storage, and loads metadata on-demand.
23///
24/// Parsed metadata is cached globally across all readers and queries via a shared
25/// `MetadataCache`, significantly reducing repeated object-storage footer reads for
26/// repeated queries on the same partitions.
27///
28/// File content caching is the responsibility of `object_store` itself: this
29/// factory is typically constructed with an object store already wrapped by the
30/// in-process L1 cache (see `object_cache::l1_wrap`), so it just reads bytes
31/// through it directly.
32pub struct ReaderFactory {
33    object_store: Arc<dyn ObjectStore>,
34    metadata_cache: Arc<MetadataCache>,
35}
36
37impl std::fmt::Debug for ReaderFactory {
38    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
39        f.debug_struct("ReaderFactory")
40            .field("metadata_cache", &self.metadata_cache)
41            .finish()
42    }
43}
44
45impl ReaderFactory {
46    /// Creates a new ReaderFactory with a shared metadata cache, reading file
47    /// contents through `object_store`.
48    pub fn new(object_store: Arc<dyn ObjectStore>, metadata_cache: Arc<MetadataCache>) -> Self {
49        Self {
50            object_store,
51            metadata_cache,
52        }
53    }
54}
55
56impl ParquetFileReaderFactory for ReaderFactory {
57    fn create_reader(
58        &self,
59        partition_index: usize,
60        partitioned_file: PartitionedFile,
61        _metadata_size_hint: Option<usize>,
62        metrics: &ExecutionPlanMetricsSet,
63    ) -> datafusion::error::Result<Box<dyn AsyncFileReader + Send>> {
64        let path = partitioned_file.path().clone();
65        let filename = path.to_string();
66        let file_size = partitioned_file.object_meta.size;
67        let file_metrics = ParquetFileMetrics::new(partition_index, path.as_ref(), metrics);
68
69        Ok(Box::new(ParquetReader {
70            filename,
71            file_size,
72            metadata_cache: Arc::clone(&self.metadata_cache),
73            object_store: Arc::clone(&self.object_store),
74            path,
75            bytes_scanned: file_metrics.bytes_scanned,
76        }))
77    }
78}
79
80/// Reads a parquet file's bytes and metadata directly from its `object_store`
81/// (which may itself be L1-cache-backed, see `object_cache::l1_wrap`) and a
82/// shared `MetadataCache`.
83pub struct ParquetReader {
84    pub filename: String,
85    pub file_size: u64,
86    pub metadata_cache: Arc<MetadataCache>,
87    pub object_store: Arc<dyn ObjectStore>,
88    pub path: Path,
89    pub bytes_scanned: Count,
90}
91
92impl AsyncFileReader for ParquetReader {
93    fn get_bytes(
94        &mut self,
95        range: Range<u64>,
96    ) -> BoxFuture<'_, datafusion::parquet::errors::Result<Bytes>> {
97        let filename = self.filename.clone();
98        let file_size = self.file_size;
99        let bytes_requested = range.end - range.start;
100        let object_store = Arc::clone(&self.object_store);
101        let path = self.path.clone();
102        let bytes_scanned = self.bytes_scanned.clone();
103
104        Box::pin(async move {
105            let start = std::time::Instant::now();
106            let result = object_store
107                .get_range(&path, range)
108                .await
109                .map_err(|e| datafusion::parquet::errors::ParquetError::External(Box::new(e)));
110            let duration_ms = start.elapsed().as_millis();
111
112            debug!(
113                "parquet_read file={filename} file_size={file_size} bytes={bytes_requested} duration_ms={duration_ms}"
114            );
115            bytes_scanned.add(bytes_requested as usize);
116
117            result
118        })
119    }
120
121    fn get_byte_ranges(
122        &mut self,
123        ranges: Vec<Range<u64>>,
124    ) -> BoxFuture<'_, datafusion::parquet::errors::Result<Vec<Bytes>>> {
125        let filename = self.filename.clone();
126        let file_size = self.file_size;
127        let num_ranges = ranges.len();
128        let total_bytes: u64 = ranges.iter().map(|r| r.end - r.start).sum();
129        let object_store = Arc::clone(&self.object_store);
130        let path = self.path.clone();
131        let bytes_scanned = self.bytes_scanned.clone();
132
133        Box::pin(async move {
134            let start = std::time::Instant::now();
135            let result = object_store
136                .get_ranges(&path, &ranges)
137                .await
138                .map_err(|e| datafusion::parquet::errors::ParquetError::External(Box::new(e)));
139            let duration_ms = start.elapsed().as_millis();
140
141            debug!(
142                "parquet_read file={filename} file_size={file_size} ranges={num_ranges} bytes={total_bytes} duration_ms={duration_ms}"
143            );
144            bytes_scanned.add(total_bytes as usize);
145
146            result
147        })
148    }
149
150    fn get_metadata(
151        &mut self,
152        _options: Option<&ArrowReaderOptions>,
153    ) -> BoxFuture<'_, datafusion::parquet::errors::Result<Arc<ParquetMetaData>>> {
154        let metadata_cache = Arc::clone(&self.metadata_cache);
155        let object_store = Arc::clone(&self.object_store);
156        let path = self.path.clone();
157        let file_size = self.file_size;
158        Box::pin(async move {
159            load_partition_metadata(&object_store, &path, file_size, Some(&metadata_cache)).await
160        })
161    }
162}