micromegas_analytics/lakehouse/
reader_factory.rs1use 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
21pub 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 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
80pub 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}