micromegas_analytics/lakehouse/
sql_batch_view.rs1use 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
31pub type MergerMaker =
33 dyn Fn(Arc<RuntimeEnv>, Arc<Schema>) -> Arc<dyn PartitionMerger> + Send + Sync;
34
35#[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 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}