Skip to main content

micromegas_analytics/lakehouse/
regenerate_partitions_table_function.rs

1use 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/// A DataFusion `TableFunctionImpl` for force-regenerating lakehouse partitions directly from
23/// source data, bypassing the "already up to date" freshness check `materialize_partitions` stops
24/// at. See `tasks/blocks_view_ordered_merges_plan.md`'s Design ยง3.
25#[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        // an alternative would be to use coerce & create_physical_expr
77        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            // Keep a clone of the raw sender alongside the LogSender wrapping it: ordinary
96            // progress lines flow through the logger as Ok((time, msg)), matching
97            // materialize_partitions's behavior, but a regenerate_partition_range failure must
98            // surface as a query-level error, not just one more log row -- so it is sent as a
99            // single Err item through the raw sender instead of only being logged.
100            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}