1use crate::time::TimeRange;
2
3use super::{partition::Partition, view::ViewMetadata};
4use anyhow::{Context, Result};
5use async_trait::async_trait;
6use chrono::{DateTime, Utc};
7use micromegas_tracing::prelude::*;
8use sqlx::{PgPool, Row};
9use std::{fmt, sync::Arc};
10
11#[async_trait]
13pub trait QueryPartitionProvider: std::fmt::Display + Send + Sync + std::fmt::Debug {
14 async fn fetch(
16 &self,
17 view_set_name: &str,
18 view_instance_id: &str,
19 query_range: Option<TimeRange>,
20 file_schema_hash: Vec<u8>,
21 ) -> Result<Vec<Partition>>;
22}
23
24#[derive(Debug)]
26pub struct PartitionCache {
27 pub partitions: Vec<Partition>,
28 insert_range: TimeRange,
29}
30
31impl fmt::Display for PartitionCache {
32 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
33 write!(f, "{self:?}")
34 }
35}
36
37impl PartitionCache {
38 pub fn empty(insert_range: TimeRange) -> Self {
42 Self {
43 partitions: vec![],
44 insert_range,
45 }
46 }
47
48 pub fn len(&self) -> usize {
49 self.partitions.len()
50 }
51
52 pub fn is_empty(&self) -> bool {
53 self.partitions.is_empty()
54 }
55
56 #[span_fn]
60 pub async fn fetch_overlapping_insert_range(
61 pool: &sqlx::PgPool,
62 insert_range: TimeRange,
63 ) -> Result<Self> {
64 let rows = instrument_named!(
65 sqlx::query(
66 "SELECT view_set_name,
67 view_instance_id,
68 begin_insert_time,
69 end_insert_time,
70 min_event_time,
71 max_event_time,
72 updated,
73 file_path,
74 file_size,
75 file_schema_hash,
76 source_data_hash,
77 num_rows,
78 sort_order
79 FROM lakehouse_partitions
80 WHERE begin_insert_time < $1
81 AND end_insert_time > $2
82 ORDER BY begin_insert_time, file_path
83 ;",
84 )
85 .bind(insert_range.end)
86 .bind(insert_range.begin)
87 .fetch_all(pool),
88 "sql_select_overlapping_partitions"
89 )
90 .await
91 .with_context(|| "fetching partitions")?;
92 let mut partitions = vec![];
93 for r in rows {
94 let view_metadata = ViewMetadata {
95 view_set_name: Arc::new(r.try_get("view_set_name")?),
96 view_instance_id: Arc::new(r.try_get("view_instance_id")?),
97 file_schema_hash: r.try_get("file_schema_hash")?,
98 };
99 let insert_time_range = TimeRange {
101 begin: r.try_get("begin_insert_time")?,
102 end: r.try_get("end_insert_time")?,
103 };
104 let event_time_range = match (
105 r.try_get::<DateTime<Utc>, _>("min_event_time").ok(),
106 r.try_get::<DateTime<Utc>, _>("max_event_time").ok(),
107 ) {
108 (Some(begin), Some(end)) => Some(TimeRange { begin, end }),
109 (None, None) => None, (Some(_), None) | (None, Some(_)) => {
111 anyhow::bail!(
112 "Corrupt partition record: only one of min/max_event_time is NULL"
113 );
114 }
115 };
116 let partition = Partition {
117 view_metadata,
118 insert_time_range,
119 event_time_range,
120 updated: r.try_get("updated")?,
121 file_path: r.try_get::<String, _>("file_path").ok(),
122 file_size: r.try_get("file_size")?,
123 source_data_hash: r.try_get("source_data_hash")?,
124 num_rows: r.try_get("num_rows")?,
125 sort_order: r.try_get("sort_order")?,
126 };
127 partition
128 .validate()
129 .with_context(|| "validating partition from database")?;
130 partitions.push(partition);
131 }
132 Ok(Self {
133 partitions,
134 insert_range,
135 })
136 }
137
138 #[span_fn]
140 pub async fn fetch_overlapping_insert_range_for_view(
141 pool: &sqlx::PgPool,
142 view_set_name: Arc<String>,
143 view_instance_id: Arc<String>,
144 insert_range: TimeRange,
145 ) -> Result<Self> {
146 let rows = instrument_named!(
147 sqlx::query(
148 "SELECT begin_insert_time,
149 end_insert_time,
150 min_event_time,
151 max_event_time,
152 updated,
153 file_path,
154 file_size,
155 file_schema_hash,
156 source_data_hash,
157 num_rows,
158 sort_order
159 FROM lakehouse_partitions
160 WHERE begin_insert_time < $1
161 AND end_insert_time > $2
162 AND view_set_name = $3
163 AND view_instance_id = $4
164 ORDER BY begin_insert_time, file_path
165 ;",
166 )
167 .bind(insert_range.end)
168 .bind(insert_range.begin)
169 .bind(&*view_set_name)
170 .bind(&*view_instance_id)
171 .fetch_all(pool),
172 "sql_select_overlapping_partitions_for_view"
173 )
174 .await
175 .with_context(|| "fetching partitions")?;
176 let mut partitions = vec![];
177 for r in rows {
178 let view_metadata = ViewMetadata {
179 view_set_name: view_set_name.clone(),
180 view_instance_id: view_instance_id.clone(),
181 file_schema_hash: r.try_get("file_schema_hash")?,
182 };
183 let insert_time_range = TimeRange {
185 begin: r.try_get("begin_insert_time")?,
186 end: r.try_get("end_insert_time")?,
187 };
188 let event_time_range = match (
189 r.try_get::<DateTime<Utc>, _>("min_event_time").ok(),
190 r.try_get::<DateTime<Utc>, _>("max_event_time").ok(),
191 ) {
192 (Some(begin), Some(end)) => Some(TimeRange { begin, end }),
193 (None, None) => None, (Some(_), None) | (None, Some(_)) => {
195 anyhow::bail!(
196 "Corrupt partition record: only one of min/max_event_time is NULL"
197 );
198 }
199 };
200 let partition = Partition {
201 view_metadata,
202 insert_time_range,
203 event_time_range,
204 updated: r.try_get("updated")?,
205 file_path: r.try_get::<String, _>("file_path").ok(),
206 file_size: r.try_get("file_size")?,
207 source_data_hash: r.try_get("source_data_hash")?,
208 num_rows: r.try_get("num_rows")?,
209 sort_order: r.try_get("sort_order")?,
210 };
211 partition
212 .validate()
213 .with_context(|| "validating partition from database")?;
214 partitions.push(partition);
215 }
216 Ok(Self {
217 partitions,
218 insert_range,
219 })
220 }
221
222 pub fn filter(
224 &self,
225 view_set_name: &str,
226 view_instance_id: &str,
227 file_schema_hash: &[u8],
228 insert_range: TimeRange,
229 ) -> Self {
230 let mut partitions = vec![];
231 for part in &self.partitions {
232 if *part.view_metadata.view_set_name == view_set_name
233 && *part.view_metadata.view_instance_id == view_instance_id
234 && part.view_metadata.file_schema_hash == file_schema_hash
235 && part.begin_insert_time() < insert_range.end
236 && part.end_insert_time() > insert_range.begin
237 {
238 partitions.push(part.clone());
239 }
240 }
241 Self {
242 partitions,
243 insert_range,
244 }
245 }
246
247 pub fn filter_insert_range(&self, insert_range: TimeRange) -> Self {
249 let mut partitions = vec![];
250 for part in &self.partitions {
251 if part.begin_insert_time() < insert_range.end
252 && part.end_insert_time() > insert_range.begin
253 {
254 partitions.push(part.clone());
255 }
256 }
257 Self {
258 partitions,
259 insert_range,
260 }
261 }
262
263 pub fn filter_inside_range(
265 &self,
266 view_set_name: &str,
267 view_instance_id: &str,
268 insert_range: TimeRange,
269 ) -> Self {
270 let mut partitions = vec![];
271 for part in &self.partitions {
272 if *part.view_metadata.view_set_name == view_set_name
273 && *part.view_metadata.view_instance_id == view_instance_id
274 && part.begin_insert_time() >= insert_range.begin
275 && part.end_insert_time() <= insert_range.end
276 {
277 partitions.push(part.clone());
278 }
279 }
280 Self {
281 partitions,
282 insert_range,
283 }
284 }
285}
286
287#[async_trait]
288impl QueryPartitionProvider for PartitionCache {
289 #[span_fn]
291 async fn fetch(
292 &self,
293 view_set_name: &str,
294 view_instance_id: &str,
295 query_range: Option<TimeRange>,
296 file_schema_hash: Vec<u8>,
297 ) -> Result<Vec<Partition>> {
298 let mut partitions = vec![];
299 if let Some(range) = query_range {
300 if range.begin < self.insert_range.begin || range.end > self.insert_range.end {
301 anyhow::bail!("filtering from a result set that's not large enough");
302 }
303 for part in &self.partitions {
304 if *part.view_metadata.view_set_name == view_set_name
305 && *part.view_metadata.view_instance_id == view_instance_id
306 && part.begin_insert_time() < range.end
307 && part.end_insert_time() > range.begin
308 && part.view_metadata.file_schema_hash == file_schema_hash
309 {
310 partitions.push(part.clone());
311 }
312 }
313 } else {
314 for part in &self.partitions {
315 if *part.view_metadata.view_set_name == view_set_name
316 && *part.view_metadata.view_instance_id == view_instance_id
317 && part.view_metadata.file_schema_hash == file_schema_hash
318 {
319 partitions.push(part.clone());
320 }
321 }
322 }
323 Ok(partitions)
324 }
325}
326
327#[derive(Debug)]
329pub struct LivePartitionProvider {
330 db_pool: PgPool,
331}
332
333impl fmt::Display for LivePartitionProvider {
334 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
335 write!(f, "{self:?}")
336 }
337}
338
339impl LivePartitionProvider {
340 pub fn new(db_pool: PgPool) -> Self {
341 Self { db_pool }
342 }
343}
344
345#[async_trait]
346impl QueryPartitionProvider for LivePartitionProvider {
347 #[span_fn]
348 async fn fetch(
349 &self,
350 view_set_name: &str,
351 view_instance_id: &str,
352 query_range: Option<TimeRange>,
353 file_schema_hash: Vec<u8>,
354 ) -> Result<Vec<Partition>> {
355 let mut partitions = vec![];
356 let rows = if let Some(range) = query_range {
357 instrument_named!(
358 sqlx::query(
359 "SELECT view_set_name,
360 view_instance_id,
361 begin_insert_time,
362 end_insert_time,
363 min_event_time,
364 max_event_time,
365 updated,
366 file_path,
367 file_size,
368 file_schema_hash,
369 source_data_hash,
370 num_rows,
371 sort_order
372 FROM lakehouse_partitions
373 WHERE view_set_name = $1
374 AND view_instance_id = $2
375 AND min_event_time <= $3
376 AND max_event_time >= $4
377 AND file_schema_hash = $5
378 ORDER BY begin_insert_time, file_path
379 ;",
380 )
381 .bind(view_set_name)
382 .bind(view_instance_id)
383 .bind(range.end)
384 .bind(range.begin)
385 .bind(file_schema_hash)
386 .fetch_all(&self.db_pool),
387 "sql_select_live_partitions"
388 )
389 .await
390 .with_context(|| "listing lakehouse partitions")?
391 } else {
392 instrument_named!(
393 sqlx::query(
394 "SELECT view_set_name,
395 view_instance_id,
396 begin_insert_time,
397 end_insert_time,
398 min_event_time,
399 max_event_time,
400 updated,
401 file_path,
402 file_size,
403 file_schema_hash,
404 source_data_hash,
405 num_rows,
406 sort_order
407 FROM lakehouse_partitions
408 WHERE view_set_name = $1
409 AND view_instance_id = $2
410 AND file_schema_hash = $3
411 ORDER BY begin_insert_time, file_path
412 ;",
413 )
414 .bind(view_set_name)
415 .bind(view_instance_id)
416 .bind(file_schema_hash)
417 .fetch_all(&self.db_pool),
418 "sql_select_live_partitions"
419 )
420 .await
421 .with_context(|| "listing lakehouse partitions")?
422 };
423 for r in rows {
424 let view_metadata = ViewMetadata {
425 view_set_name: Arc::new(r.try_get("view_set_name")?),
426 view_instance_id: Arc::new(r.try_get("view_instance_id")?),
427 file_schema_hash: r.try_get("file_schema_hash")?,
428 };
429 let insert_time_range = TimeRange {
431 begin: r.try_get("begin_insert_time")?,
432 end: r.try_get("end_insert_time")?,
433 };
434 let event_time_range = match (
435 r.try_get::<DateTime<Utc>, _>("min_event_time").ok(),
436 r.try_get::<DateTime<Utc>, _>("max_event_time").ok(),
437 ) {
438 (Some(begin), Some(end)) => Some(TimeRange { begin, end }),
439 (None, None) => None, (Some(_), None) | (None, Some(_)) => {
441 anyhow::bail!(
442 "Corrupt partition record: only one of min/max_event_time is NULL"
443 );
444 }
445 };
446 let partition = Partition {
447 view_metadata,
448 insert_time_range,
449 event_time_range,
450 updated: r.try_get("updated")?,
451 file_path: r.try_get::<String, _>("file_path").ok(),
452 file_size: r.try_get("file_size")?,
453 source_data_hash: r.try_get("source_data_hash")?,
454 num_rows: r.try_get("num_rows")?,
455 sort_order: r.try_get("sort_order")?,
456 };
457 partition
458 .validate()
459 .with_context(|| "validating partition from database")?;
460 partitions.push(partition);
461 }
462 Ok(partitions)
463 }
464}
465
466#[derive(Debug)]
468pub struct NullPartitionProvider {}
469
470impl fmt::Display for NullPartitionProvider {
471 fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
472 write!(f, "{self:?}")
473 }
474}
475
476#[async_trait]
477impl QueryPartitionProvider for NullPartitionProvider {
478 async fn fetch(
479 &self,
480 _view_set_name: &str,
481 _view_instance_id: &str,
482 _query_range: Option<TimeRange>,
483 _file_schema_hash: Vec<u8>,
484 ) -> Result<Vec<Partition>> {
485 Ok(vec![])
486 }
487}