Skip to main content

micromegas_analytics/lakehouse/
sql_batch_view.rs

1use super::{
2    batch_update::PartitionCreationStrategy,
3    dataframe_time_bounds::{DataFrameTimeBounds, NamedColumnsTimeBounds},
4    lakehouse_context::LakehouseContext,
5    materialized_view::MaterializedView,
6    merge::{MergeQueryResult, PartitionMerger, QueryMerger},
7    partition::Partition,
8    partition_cache::{NullPartitionProvider, PartitionCache},
9    query::make_session_context,
10    session_configurator::SessionConfigurator,
11    sql_partition_spec::fetch_sql_partition_spec,
12    view::{PartitionSpec, View, ViewMetadata},
13    view_factory::ViewFactory,
14};
15use crate::{
16    record_batch_transformer::TrivialRecordBatchTransformer,
17    time::{TimeRange, datetime_to_scalar},
18};
19use anyhow::{Context, Result};
20use async_trait::async_trait;
21use chrono::{DateTime, TimeDelta, Utc};
22use datafusion::{
23    arrow::datatypes::Schema, execution::runtime_env::RuntimeEnv, prelude::*, sql::TableReference,
24};
25use micromegas_ingestion::data_lake_connection::DataLakeConnection;
26use micromegas_tracing::error;
27use std::hash::Hash;
28use std::hash::Hasher;
29use std::{hash::DefaultHasher, sync::Arc};
30
31/// A type alias for a function that creates a `PartitionMerger`.
32pub type MergerMaker =
33    dyn Fn(Arc<RuntimeEnv>, Arc<Schema>) -> Arc<dyn PartitionMerger> + Send + Sync;
34
35/// SQL-defined view updated in batch
36#[derive(Debug)]
37pub struct SqlBatchView {
38    view_set_name: Arc<String>,
39    view_instance_id: Arc<String>,
40    min_event_time_column: Arc<String>,
41    max_event_time_column: Arc<String>,
42    count_src_query: Arc<String>,
43    extract_query: Arc<String>,
44    merge_partitions_query: Arc<String>,
45    schema: Arc<Schema>,
46    merger: Arc<dyn PartitionMerger>,
47    view_factory: Arc<ViewFactory>,
48    session_configurator: Arc<dyn SessionConfigurator>,
49    update_group: Option<i32>,
50    max_partition_delta_from_source: TimeDelta,
51    max_partition_delta_from_merge: TimeDelta,
52}
53
54impl SqlBatchView {
55    #[expect(clippy::too_many_arguments)]
56    /// # Arguments
57    ///
58    /// * `runtime` - datafusion runtime
59    /// * `view_set_name` - name of the table
60    /// * `min_event_time_column` - min(column) should result in the first timestamp in a dataframe
61    /// * `max_event_time_column` - max(column) should result in the last timestamp in a dataframe
62    /// * `count_src_query` - used to count the rows of the underlying data to know if a cached partition is up to date
63    /// * `extract_query` - used to extract the source data into a cached partition
64    /// * `merge_partitions_query` - used to merge multiple partitions into a single one (and user queries which are one multiple partitions by default)
65    /// * `lake` - data lake
66    /// * `view_factory` - all views accessible to the `count_src_query`
67    /// * `session_configurator` - configurator for registering custom tables (e.g., JSON files)
68    /// * `update_group` - tells the daemon which view should be materialized and in what order
69    pub async fn new(
70        runtime: Arc<RuntimeEnv>,
71        view_set_name: Arc<String>,
72        min_event_time_column: Arc<String>,
73        max_event_time_column: Arc<String>,
74        count_src_query: Arc<String>,
75        extract_query: Arc<String>,
76        merge_partitions_query: Arc<String>,
77        lake: Arc<DataLakeConnection>,
78        view_factory: Arc<ViewFactory>,
79        session_configurator: Arc<dyn SessionConfigurator>,
80        update_group: Option<i32>,
81        max_partition_delta_from_source: TimeDelta,
82        max_partition_delta_from_merge: TimeDelta,
83        merger_maker: Option<&MergerMaker>,
84    ) -> Result<Self> {
85        let null_part_provider = Arc::new(NullPartitionProvider {});
86        let lakehouse = Arc::new(LakehouseContext::new(lake.clone(), runtime.clone()));
87        let ctx = make_session_context(
88            lakehouse,
89            null_part_provider,
90            None,
91            view_factory.clone(),
92            session_configurator.clone(),
93            true,
94        )
95        .await
96        .with_context(|| "make_session_context")?;
97        let now_str = Utc::now().to_rfc3339();
98        let sql = extract_query
99            .replace("{begin}", &now_str)
100            .replace("{end}", &now_str);
101        let extracted_df = ctx.sql(&sql).await?;
102        let schema = extracted_df.schema().inner().clone();
103        let session_configurator_for_merger = session_configurator.clone();
104        let merger = merger_maker.unwrap_or(&|_runtime, schema| {
105            let merge_query = Arc::new(merge_partitions_query.replace("{source}", "source"));
106            Arc::new(QueryMerger::new(
107                view_factory.clone(),
108                session_configurator_for_merger.clone(),
109                schema,
110                merge_query,
111            ))
112        })(runtime.clone(), schema.clone());
113
114        Ok(Self {
115            view_set_name,
116            view_instance_id: Arc::new(String::from("global")),
117            min_event_time_column,
118            max_event_time_column,
119            count_src_query,
120            extract_query,
121            merge_partitions_query,
122            schema,
123            merger,
124            view_factory,
125            session_configurator,
126            update_group,
127            max_partition_delta_from_source,
128            max_partition_delta_from_merge,
129        })
130    }
131}
132
133#[async_trait]
134impl View for SqlBatchView {
135    fn get_view_set_name(&self) -> Arc<String> {
136        self.view_set_name.clone()
137    }
138
139    fn get_view_instance_id(&self) -> Arc<String> {
140        self.view_instance_id.clone()
141    }
142
143    async fn make_batch_partition_spec(
144        &self,
145        lakehouse: Arc<LakehouseContext>,
146        existing_partitions: Arc<PartitionCache>,
147        insert_range: TimeRange,
148    ) -> Result<Arc<dyn PartitionSpec>> {
149        let view_meta = ViewMetadata {
150            view_set_name: self.get_view_set_name(),
151            view_instance_id: self.get_view_instance_id(),
152            file_schema_hash: self.get_file_schema_hash(),
153        };
154        let partitions_in_range = Arc::new(existing_partitions.filter_insert_range(insert_range));
155        let ctx = make_session_context(
156            lakehouse,
157            partitions_in_range.clone(),
158            None,
159            self.view_factory.clone(),
160            self.session_configurator.clone(),
161            true,
162        )
163        .await
164        .with_context(|| "make_session_context")?;
165
166        let count_src_sql = self
167            .count_src_query
168            .replace("{begin}", &insert_range.begin.to_rfc3339())
169            .replace("{end}", &insert_range.end.to_rfc3339());
170
171        let extract_sql = self
172            .extract_query
173            .replace("{begin}", &insert_range.begin.to_rfc3339())
174            .replace("{end}", &insert_range.end.to_rfc3339());
175
176        Ok(Arc::new(
177            fetch_sql_partition_spec(
178                ctx,
179                Arc::new(TrivialRecordBatchTransformer {}),
180                self.get_time_bounds(),
181                self.schema.clone(),
182                count_src_sql,
183                extract_sql,
184                view_meta,
185                insert_range,
186            )
187            .await
188            .with_context(|| "fetch_sql_partition_spec")?,
189        ))
190    }
191
192    fn get_file_schema_hash(&self) -> Vec<u8> {
193        let mut hasher = DefaultHasher::new();
194        self.schema.hash(&mut hasher);
195        hasher.finish().to_le_bytes().to_vec()
196    }
197
198    fn get_file_schema(&self) -> Arc<Schema> {
199        self.schema.clone()
200    }
201
202    async fn jit_update(
203        &self,
204        _lakehouse: Arc<LakehouseContext>,
205        _query_range: Option<TimeRange>,
206    ) -> Result<()> {
207        Ok(())
208    }
209
210    fn make_time_filter(&self, begin: DateTime<Utc>, end: DateTime<Utc>) -> Result<Vec<Expr>> {
211        Ok(vec![
212            col(&*self.min_event_time_column).lt_eq(lit(datetime_to_scalar(end))),
213            col(&*self.max_event_time_column).gt_eq(lit(datetime_to_scalar(begin))),
214        ])
215    }
216
217    fn get_time_bounds(&self) -> Arc<dyn DataFrameTimeBounds> {
218        Arc::new(NamedColumnsTimeBounds::new(
219            self.min_event_time_column.clone(),
220            self.max_event_time_column.clone(),
221        ))
222    }
223
224    async fn register_table(&self, ctx: &SessionContext, table: MaterializedView) -> Result<()> {
225        let view_name = self.get_view_set_name().to_string();
226        let partitions_table_name = format!("__{view_name}__partitions");
227        ctx.register_table(
228            TableReference::Bare {
229                table: partitions_table_name.clone().into(),
230            },
231            Arc::new(table),
232        )?;
233        let df = ctx
234            .sql(
235                &self
236                    .merge_partitions_query
237                    .replace("{source}", &partitions_table_name),
238            )
239            .await?;
240        ctx.register_table(
241            TableReference::Bare {
242                table: view_name.into(),
243            },
244            df.into_view(),
245        )?;
246        Ok(())
247    }
248
249    async fn merge_partitions(
250        &self,
251        lakehouse: Arc<LakehouseContext>,
252        partitions_to_merge: Arc<Vec<Partition>>,
253        partitions_all_views: Arc<PartitionCache>,
254        insert_range: TimeRange,
255    ) -> Result<MergeQueryResult> {
256        let res = self
257            .merger
258            .execute_merge_query(
259                lakehouse,
260                partitions_to_merge,
261                partitions_all_views,
262                insert_range,
263            )
264            .await;
265        if let Err(e) = &res {
266            error!("{e:?}");
267        }
268        res
269    }
270
271    fn get_update_group(&self) -> Option<i32> {
272        self.update_group
273    }
274
275    fn get_max_partition_time_delta(&self, strategy: &PartitionCreationStrategy) -> TimeDelta {
276        match strategy {
277            PartitionCreationStrategy::Abort | PartitionCreationStrategy::CreateFromSource => {
278                self.max_partition_delta_from_source
279            }
280            PartitionCreationStrategy::MergeExisting(_partitions) => {
281                self.max_partition_delta_from_merge
282            }
283        }
284    }
285}