micromegas_analytics/lakehouse/
regenerate_partitions_table_function.rs1use super::batch_update::regenerate_partition_range;
2use super::lakehouse_context::LakehouseContext;
3use super::partition_cache::PartitionCache;
4use super::view_factory::ViewFactory;
5use crate::dfext::expressions::exp_to_i64;
6use crate::dfext::expressions::exp_to_string;
7use crate::dfext::expressions::exp_to_timestamp;
8use crate::dfext::log_stream_table_provider::LogStreamTableProvider;
9use crate::dfext::task_log_exec_plan::TaskLogExecPlan;
10use crate::response_writer::LogSender;
11use crate::response_writer::Logger;
12use crate::time::TimeRange;
13use anyhow::Context;
14use chrono::TimeDelta;
15use datafusion::catalog::TableFunctionArgs;
16use datafusion::catalog::TableFunctionImpl;
17use datafusion::catalog::TableProvider;
18use datafusion::common::plan_err;
19use micromegas_tracing::prelude::*;
20use std::sync::Arc;
21
22#[derive(Debug)]
26pub struct RegeneratePartitionsTableFunction {
27 lakehouse: Arc<LakehouseContext>,
28 view_factory: Arc<ViewFactory>,
29}
30
31impl RegeneratePartitionsTableFunction {
32 pub fn new(lakehouse: Arc<LakehouseContext>, view_factory: Arc<ViewFactory>) -> Self {
33 Self {
34 lakehouse,
35 view_factory,
36 }
37 }
38}
39
40#[span_fn]
41async fn regenerate_partitions_impl(
42 lakehouse: Arc<LakehouseContext>,
43 view_factory: Arc<ViewFactory>,
44 view_name: &str,
45 insert_range: TimeRange,
46 partition_time_delta: TimeDelta,
47 logger: Arc<dyn Logger>,
48) -> anyhow::Result<()> {
49 let view = view_factory
50 .get_global_view(view_name)
51 .with_context(|| format!("can't find view {view_name}"))?;
52
53 let existing_partitions_all_views = Arc::new(
54 PartitionCache::fetch_overlapping_insert_range(&lakehouse.lake().db_pool, insert_range)
55 .await?,
56 );
57
58 regenerate_partition_range(
59 existing_partitions_all_views,
60 lakehouse,
61 view,
62 insert_range,
63 partition_time_delta,
64 logger,
65 )
66 .await?;
67 Ok(())
68}
69
70impl TableFunctionImpl for RegeneratePartitionsTableFunction {
71 fn call_with_args(
72 &self,
73 args: TableFunctionArgs,
74 ) -> datafusion::error::Result<Arc<dyn TableProvider>> {
75 let args = args.exprs();
76 let Some(view_set_name) = args.first().map(exp_to_string).transpose()? else {
78 return plan_err!("Missing first argument, expected view_set_name: String");
79 };
80 let Some(begin) = args.get(1).map(exp_to_timestamp).transpose()? else {
81 return plan_err!("Missing 2nd argument, expected a UTC nanoseconds timestamp");
82 };
83 let Some(end) = args.get(2).map(exp_to_timestamp).transpose()? else {
84 return plan_err!("Missing 3rd argument, expected a UTC nanoseconds timestamp");
85 };
86 let Some(delta) = args.get(3).map(exp_to_i64).transpose()? else {
87 return plan_err!("Missing 4th argument, expected a number of seconds(i64)");
88 };
89
90 let lakehouse = self.lakehouse.clone();
91 let view_factory = self.view_factory.clone();
92
93 let spawner = move || {
94 let (tx, rx) = tokio::sync::mpsc::channel(100);
95 let raw_tx = tx.clone();
101 let logger = Arc::new(LogSender::new(tx));
102 spawn_with_context(async move {
103 if let Err(e) = regenerate_partitions_impl(
104 lakehouse,
105 view_factory,
106 &view_set_name,
107 TimeRange::new(begin, end),
108 TimeDelta::seconds(delta),
109 logger.clone(),
110 )
111 .await
112 .with_context(|| "regenerate_partitions_impl")
113 {
114 let msg = format!("{e:?}");
115 error!("{msg}");
116 let _ = raw_tx.send(Err(msg)).await;
117 }
118 });
119 rx
120 };
121
122 Ok(Arc::new(LogStreamTableProvider {
123 log_stream: Arc::new(TaskLogExecPlan::new(Box::new(spawner))),
124 }))
125 }
126}