micromegas_analytics/lakehouse/
materialized_view.rs1use 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#[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 fn supports_filters_pushdown(
103 &self,
104 filters: &[&Expr],
105 ) -> datafusion::error::Result<Vec<TableProviderFilterPushDown>> {
106 Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
110 }
111}