Skip to main content

micromegas_analytics/lakehouse/
partition_cache.rs

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/// A trait for providing queryable partitions.
12#[async_trait]
13pub trait QueryPartitionProvider: std::fmt::Display + Send + Sync + std::fmt::Debug {
14    /// Fetches partitions based on the provided criteria.
15    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/// PartitionCache allows to query partitions based on the insert_time range
25#[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    /// Builds an empty cache for `insert_range`, without querying Postgres. Used by
39    /// offline/no-DB tests that need a `PartitionCache` handle to satisfy an API (e.g. session
40    /// context registration) but have no other view's partitions to supply.
41    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    /// fetches the partitions of all views matching the specified insert range
57    //todo: this should be limited to global instances
58    //todo: ask for a list of view sets (which would be provided by the views using a get_dependencies() api entry)
59    #[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            // file_metadata will be loaded on-demand when needed
100            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, // Empty partition - both NULL
110                (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    /// fetches the partitions of a single view instance matching the specified insert range
139    #[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            // file_metadata will be loaded on-demand when needed
184            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, // Empty partition - both NULL
194                (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    // overlap test for a specific view
223    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    // overlap test for a all views
248    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    // single view that fits completely in the specified range
264    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    /// unlike LivePartitionProvider, the query_range is tested against the insertion time, not the event time
290    #[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/// A `QueryPartitionProvider` that fetches partitions directly from the database.
328#[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            // file_metadata will be loaded on-demand when needed
430            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, // Empty partition - both NULL
440                (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/// A `QueryPartitionProvider` that always returns an empty list of partitions.
467#[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}