Skip to main content

micromegas_analytics/lakehouse/
view_instance_table_function.rs

1use super::{
2    lakehouse_context::LakehouseContext, materialized_view::MaterializedView,
3    partition_cache::QueryPartitionProvider, view_factory::ViewFactory,
4};
5use crate::{dfext::expressions::exp_to_string, time::TimeRange};
6use datafusion::{
7    catalog::{TableFunctionArgs, TableFunctionImpl, TableProvider},
8    common::plan_err,
9    error::DataFusionError,
10};
11use micromegas_tracing::prelude::*;
12use std::sync::Arc;
13
14/// `ViewInstanceTableFunction` gives access to any view instance using a [ViewFactory].
15///
16/// ```python
17/// # Python code showing the usage of `view_instance(view_set_name, view_instance_id)`
18/// sql = """
19/// SELECT *
20/// FROM view_instance('thread_spans', '{stream_id}')
21/// ;""".format(stream_id=stream_id)
22/// df_spans = client.query(sql, begin_spans, end_spans)
23/// ```
24///
25#[derive(Debug)]
26pub struct ViewInstanceTableFunction {
27    lakehouse: Arc<LakehouseContext>,
28    view_factory: Arc<ViewFactory>,
29    part_provider: Arc<dyn QueryPartitionProvider>,
30    query_range: Option<TimeRange>,
31}
32
33impl ViewInstanceTableFunction {
34    pub fn new(
35        lakehouse: Arc<LakehouseContext>,
36        view_factory: Arc<ViewFactory>,
37        part_provider: Arc<dyn QueryPartitionProvider>,
38        query_range: Option<TimeRange>,
39    ) -> Self {
40        Self {
41            lakehouse,
42            view_factory,
43            part_provider,
44            query_range,
45        }
46    }
47}
48
49impl TableFunctionImpl for ViewInstanceTableFunction {
50    #[span_fn]
51    fn call_with_args(
52        &self,
53        args: TableFunctionArgs,
54    ) -> datafusion::error::Result<Arc<dyn TableProvider>> {
55        let exprs = args.exprs();
56        let arg1 = exprs.first().map(exp_to_string);
57        let Some(Ok(view_set_name)) = arg1 else {
58            return plan_err!(
59                "First argument to view_instance must be a string (the view set name), given {:?}",
60                arg1
61            );
62        };
63        let arg2 = exprs.get(1).map(exp_to_string);
64        let Some(Ok(view_instance_id)) = arg2 else {
65            return plan_err!(
66                "Second argument to view_instance must be a string (the view instance id), given {:?}",
67                arg2
68            );
69        };
70
71        let view = self
72            .view_factory
73            .make_view(&view_set_name, &view_instance_id)
74            .map_err(|e| DataFusionError::Plan(format!("error making view {e:?}")))?;
75
76        Ok(Arc::new(MaterializedView::new(
77            self.lakehouse.clone(),
78            self.lakehouse.reader_factory().clone(),
79            view,
80            self.part_provider.clone(),
81            self.query_range,
82        )))
83    }
84}