Skip to main content

micromegas_analytics/lakehouse/
materialized_view.rs

1use super::{
2    lakehouse_context::LakehouseContext,
3    partition_cache::QueryPartitionProvider,
4    partitioned_execution_plan::{OrderingBounds, make_partitioned_execution_plan},
5    reader_factory::ReaderFactory,
6    view::View,
7};
8use crate::time::TimeRange;
9use async_trait::async_trait;
10use datafusion::{
11    arrow::datatypes::SchemaRef,
12    catalog::{Session, TableProvider},
13    datasource::TableType,
14    error::DataFusionError,
15    logical_expr::{Expr, TableProviderFilterPushDown},
16    physical_plan::ExecutionPlan,
17};
18use micromegas_tracing::prelude::*;
19use std::sync::Arc;
20
21/// A DataFusion `TableProvider` for materialized views.
22#[derive(Debug)]
23pub struct MaterializedView {
24    lakehouse: Arc<LakehouseContext>,
25    reader_factory: Arc<ReaderFactory>,
26    view: Arc<dyn View>,
27    part_provider: Arc<dyn QueryPartitionProvider>,
28    query_range: Option<TimeRange>,
29}
30
31impl MaterializedView {
32    pub fn new(
33        lakehouse: Arc<LakehouseContext>,
34        reader_factory: Arc<ReaderFactory>,
35        view: Arc<dyn View>,
36        part_provider: Arc<dyn QueryPartitionProvider>,
37        query_range: Option<TimeRange>,
38    ) -> Self {
39        Self {
40            lakehouse,
41            reader_factory,
42            view,
43            part_provider,
44            query_range,
45        }
46    }
47
48    pub fn get_view(&self) -> Arc<dyn View> {
49        self.view.clone()
50    }
51}
52
53#[async_trait]
54impl TableProvider for MaterializedView {
55    fn schema(&self) -> SchemaRef {
56        self.view.get_file_schema()
57    }
58
59    fn table_type(&self) -> TableType {
60        TableType::Base
61    }
62
63    #[span_fn]
64    async fn scan(
65        &self,
66        state: &dyn Session,
67        projection: Option<&Vec<usize>>,
68        filters: &[Expr],
69        limit: Option<usize>,
70    ) -> datafusion::error::Result<Arc<dyn ExecutionPlan>> {
71        self.view
72            .jit_update(self.lakehouse.clone(), self.query_range)
73            .await
74            .map_err(|e| DataFusionError::External(format!("{e:#}").into()))?;
75
76        let partitions = self
77            .part_provider
78            .fetch(
79                &self.view.get_view_set_name(),
80                &self.view.get_view_instance_id(),
81                self.query_range,
82                self.view.get_file_schema_hash(),
83            )
84            .await
85            .map_err(|e| datafusion::error::DataFusionError::External(e.into()))?;
86        trace!("MaterializedView::scan nb_partitions={}", partitions.len());
87
88        make_partitioned_execution_plan(
89            self.schema(),
90            self.reader_factory.clone(),
91            state,
92            projection,
93            filters,
94            limit,
95            Arc::new(partitions),
96            &self.view.get_scan_output_ordering(),
97            OrderingBounds::EventTime,
98        )
99    }
100
101    /// Tell DataFusion to push filters down to the scan method
102    fn supports_filters_pushdown(
103        &self,
104        filters: &[&Expr],
105    ) -> datafusion::error::Result<Vec<TableProviderFilterPushDown>> {
106        // Inexact because the pruning can't handle all expressions and pruning
107        // is not done at the row level -- there may be rows in returned files
108        // that do not pass the filter
109        Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
110    }
111}