Skip to main content

micromegas_analytics/lakehouse/
retire_partitions_table_function.rs

1use anyhow::Context;
2use chrono::DateTime;
3use chrono::Utc;
4use datafusion::catalog::TableFunctionArgs;
5use datafusion::catalog::TableFunctionImpl;
6use datafusion::catalog::TableProvider;
7use datafusion::common::plan_err;
8use micromegas_ingestion::data_lake_connection::DataLakeConnection;
9use micromegas_tracing::prelude::*;
10use std::sync::Arc;
11
12use crate::dfext::expressions::exp_to_string;
13use crate::dfext::expressions::exp_to_timestamp;
14use crate::dfext::log_stream_table_provider::LogStreamTableProvider;
15use crate::dfext::task_log_exec_plan::TaskLogExecPlan;
16use crate::response_writer::LogSender;
17use crate::response_writer::Logger;
18
19use super::write_partition::retire_partitions;
20
21/// A DataFusion `TableFunctionImpl` for retiring lakehouse partitions.
22#[derive(Debug)]
23pub struct RetirePartitionsTableFunction {
24    lake: Arc<DataLakeConnection>,
25}
26
27impl RetirePartitionsTableFunction {
28    pub fn new(lake: Arc<DataLakeConnection>) -> Self {
29        Self { lake }
30    }
31}
32
33async fn retire_partitions_impl(
34    lake: Arc<DataLakeConnection>,
35    view_set_name: &str,
36    view_instance_id: &str,
37    begin_insert_time: DateTime<Utc>,
38    end_insert_time: DateTime<Utc>,
39    logger: Arc<dyn Logger>,
40) -> anyhow::Result<()> {
41    let mut tr = lake.db_pool.begin().await?;
42    retire_partitions(
43        &mut tr,
44        view_set_name,
45        view_instance_id,
46        begin_insert_time,
47        end_insert_time,
48        logger,
49    )
50    .await?;
51    tr.commit().await.with_context(|| "commit")?;
52    Ok(())
53}
54
55impl TableFunctionImpl for RetirePartitionsTableFunction {
56    fn call_with_args(
57        &self,
58        args: TableFunctionArgs,
59    ) -> datafusion::error::Result<Arc<dyn TableProvider>> {
60        let args = args.exprs();
61        // an alternative would be to use coerce & create_physical_expr
62        let Some(view_set_name) = args.first().map(exp_to_string).transpose()? else {
63            return plan_err!("Missing first argument, expected view_set_name: String");
64        };
65        let Some(view_instance_id) = args.get(1).map(exp_to_string).transpose()? else {
66            return plan_err!("Missing 2nd argument, expected view_instance_id: String");
67        };
68        let Some(begin) = args.get(2).map(exp_to_timestamp).transpose()? else {
69            return plan_err!("Missing 3rd argument, expected a UTC nanoseconds timestamp");
70        };
71        let Some(end) = args.get(3).map(exp_to_timestamp).transpose()? else {
72            return plan_err!("Missing 4th argument, expected a UTC nanoseconds timestamp");
73        };
74
75        let lake = self.lake.clone();
76
77        let spawner = move || {
78            let (tx, rx) = tokio::sync::mpsc::channel(100);
79            let logger = Arc::new(LogSender::new(tx));
80            spawn_with_context(async move {
81                if let Err(e) = retire_partitions_impl(
82                    lake,
83                    &view_set_name,
84                    &view_instance_id,
85                    begin,
86                    end,
87                    logger,
88                )
89                .await
90                .with_context(|| "retire_partitions_impl")
91                {
92                    error!("{e:?}");
93                }
94            });
95            rx
96        };
97
98        Ok(Arc::new(LogStreamTableProvider {
99            log_stream: Arc::new(TaskLogExecPlan::new(Box::new(spawner))),
100        }))
101    }
102}