Skip to main content

micromegas_analytics/lakehouse/
query.rs

1use super::{
2    answer::Answer, get_payload_function::GetPayload, lakehouse_context::LakehouseContext,
3    list_partitions_table_function::ListPartitionsTableFunction,
4    list_view_sets_table_function::ListViewSetsTableFunction,
5    materialize_partitions_table_function::MaterializePartitionsTableFunction,
6    parse_block_table_function::ParseBlockTableFunction, partition::Partition,
7    partition_cache::QueryPartitionProvider, partitioned_table_provider::PartitionedTableProvider,
8    perfetto_trace_table_function::PerfettoTraceTableFunction,
9    process_spans_table_function::ProcessSpansTableFunction, reader_factory::ReaderFactory,
10    regenerate_partitions_table_function::RegeneratePartitionsTableFunction,
11    retire_partition_by_file_udf::make_retire_partition_by_file_udf,
12    retire_partition_by_metadata_udf::make_retire_partition_by_metadata_udf,
13    retire_partitions_table_function::RetirePartitionsTableFunction,
14    session_configurator::SessionConfigurator, view::View, view_factory::ViewFactory,
15};
16use crate::{
17    lakehouse::{
18        materialized_view::MaterializedView, table_scan_rewrite::TableScanRewrite,
19        view_instance_table_function::ViewInstanceTableFunction,
20    },
21    properties::{
22        properties_to_dict_udf::PropertiesToDict, properties_to_jsonb_udf::PropertiesToJsonb,
23    },
24    time::TimeRange,
25};
26use anyhow::{Context, Result};
27use datafusion::{
28    arrow::{array::RecordBatch, datatypes::SchemaRef},
29    execution::{context::SessionContext, object_store::ObjectStoreUrl, runtime_env::RuntimeEnv},
30    logical_expr::{ScalarUDF, async_udf::AsyncScalarUDF},
31    prelude::*,
32    sql::TableReference,
33};
34use micromegas_tracing::prelude::*;
35use std::sync::Arc;
36
37#[span_fn]
38async fn register_table(
39    lakehouse: Arc<LakehouseContext>,
40    reader_factory: Arc<ReaderFactory>,
41    part_provider: Arc<dyn QueryPartitionProvider>,
42    query_range: Option<TimeRange>,
43    ctx: &SessionContext,
44    view: Arc<dyn View>,
45) -> Result<()> {
46    let table = MaterializedView::new(
47        lakehouse,
48        reader_factory,
49        view.clone(),
50        part_provider,
51        query_range,
52    );
53    view.register_table(ctx, table).await
54}
55
56/// query_partitions_context returns a context to run queries using the partitions as the "source" table
57#[span_fn]
58pub async fn query_partitions_context(
59    runtime: Arc<RuntimeEnv>,
60    reader_factory: Arc<ReaderFactory>,
61    object_store: Arc<dyn object_store::ObjectStore>,
62    schema: SchemaRef,
63    partitions: Arc<Vec<Partition>>,
64) -> Result<SessionContext> {
65    let table = PartitionedTableProvider::new(schema, reader_factory, partitions);
66    let object_store_url = ObjectStoreUrl::parse("obj://lakehouse/").unwrap();
67    let ctx = SessionContext::new_with_config_rt(SessionConfig::default(), runtime);
68    ctx.register_object_store(object_store_url.as_ref(), object_store);
69    ctx.register_table(
70        TableReference::Bare {
71            table: "source".into(),
72        },
73        Arc::new(table),
74    )?;
75    register_extension_functions(&ctx);
76    Ok(ctx)
77}
78
79// query_partitions returns a dataframe, leaving the option of streaming the results
80#[span_fn]
81pub async fn query_partitions(
82    runtime: Arc<RuntimeEnv>,
83    reader_factory: Arc<ReaderFactory>,
84    object_store: Arc<dyn object_store::ObjectStore>,
85    schema: SchemaRef,
86    partitions: Arc<Vec<Partition>>,
87    sql: &str,
88) -> Result<DataFrame> {
89    let ctx =
90        query_partitions_context(runtime, reader_factory, object_store, schema, partitions).await?;
91    Ok(ctx.sql(sql).await?)
92}
93
94/// register functions that are part of the lakehouse architecture
95#[span_fn]
96pub fn register_lakehouse_functions(
97    ctx: &SessionContext,
98    lakehouse: Arc<LakehouseContext>,
99    part_provider: Arc<dyn QueryPartitionProvider>,
100    query_range: Option<TimeRange>,
101    view_factory: Arc<ViewFactory>,
102    is_admin: bool,
103) {
104    ctx.register_udtf(
105        "view_instance",
106        Arc::new(ViewInstanceTableFunction::new(
107            lakehouse.clone(),
108            view_factory.clone(),
109            part_provider.clone(),
110            query_range,
111        )),
112    );
113    ctx.register_udtf(
114        "list_partitions",
115        Arc::new(ListPartitionsTableFunction::new(lakehouse.lake().clone())),
116    );
117    ctx.register_udtf(
118        "list_view_sets",
119        Arc::new(ListViewSetsTableFunction::new(view_factory.clone())),
120    );
121    ctx.register_udtf(
122        "perfetto_trace_chunks",
123        Arc::new(PerfettoTraceTableFunction::new(
124            lakehouse.clone(),
125            view_factory.clone(),
126            part_provider.clone(),
127        )),
128    );
129    ctx.register_udtf(
130        "parse_block",
131        Arc::new(ParseBlockTableFunction::new(
132            lakehouse.clone(),
133            view_factory.clone(),
134            part_provider.clone(),
135            query_range,
136        )),
137    );
138    ctx.register_udtf(
139        "process_spans",
140        Arc::new(ProcessSpansTableFunction::new(
141            lakehouse.clone(),
142            view_factory.clone(),
143            part_provider.clone(),
144            query_range,
145        )),
146    );
147    ctx.register_udf(
148        AsyncScalarUDF::new(Arc::new(GetPayload::new(lakehouse.lake().clone()))).into_scalar_udf(),
149    );
150    if is_admin {
151        ctx.register_udtf(
152            "retire_partitions",
153            Arc::new(RetirePartitionsTableFunction::new(lakehouse.lake().clone())),
154        );
155        ctx.register_udtf(
156            "materialize_partitions",
157            Arc::new(MaterializePartitionsTableFunction::new(
158                lakehouse.clone(),
159                view_factory.clone(),
160            )),
161        );
162        ctx.register_udtf(
163            "regenerate_partitions",
164            Arc::new(RegeneratePartitionsTableFunction::new(
165                lakehouse.clone(),
166                view_factory.clone(),
167            )),
168        );
169        ctx.register_udf(
170            make_retire_partition_by_file_udf(lakehouse.lake().clone()).into_scalar_udf(),
171        );
172        ctx.register_udf(
173            make_retire_partition_by_metadata_udf(lakehouse.lake().clone()).into_scalar_udf(),
174        );
175    }
176}
177
178/// register functions that are not depended on the lakehouse architecture
179#[span_fn]
180pub fn register_extension_functions(ctx: &SessionContext) {
181    ctx.register_udf(ScalarUDF::from(PropertiesToDict::new()));
182    ctx.register_udf(ScalarUDF::from(PropertiesToJsonb::new()));
183    micromegas_datafusion_extensions::register_extension_udfs(ctx);
184}
185
186#[span_fn]
187pub fn register_functions(
188    ctx: &SessionContext,
189    lakehouse: Arc<LakehouseContext>,
190    part_provider: Arc<dyn QueryPartitionProvider>,
191    query_range: Option<TimeRange>,
192    view_factory: Arc<ViewFactory>,
193    is_admin: bool,
194) {
195    register_lakehouse_functions(
196        ctx,
197        lakehouse,
198        part_provider,
199        query_range,
200        view_factory,
201        is_admin,
202    );
203    register_extension_functions(ctx);
204}
205
206#[span_fn]
207pub async fn make_session_context(
208    lakehouse: Arc<LakehouseContext>,
209    part_provider: Arc<dyn QueryPartitionProvider>,
210    query_range: Option<TimeRange>,
211    view_factory: Arc<ViewFactory>,
212    configurator: Arc<dyn SessionConfigurator>,
213    is_admin: bool,
214) -> Result<SessionContext> {
215    // Disable page index reading for backward compatibility with legacy Parquet files
216    // Legacy files may have incomplete ColumnIndex metadata (missing null_pages field)
217    // which causes errors in DataFusion 51+ with Arrow 57.0 when reading page indexes
218    let config = SessionConfig::default()
219        .set_bool("datafusion.execution.parquet.enable_page_index", false)
220        .with_information_schema(true);
221    let ctx = SessionContext::new_with_config_rt(config, lakehouse.runtime().clone());
222    if let Some(range) = &query_range {
223        ctx.add_analyzer_rule(Arc::new(TableScanRewrite::new(*range)));
224    }
225    let object_store_url = ObjectStoreUrl::parse("obj://lakehouse/").unwrap();
226    let object_store = lakehouse.lake().blob_storage.inner();
227    ctx.register_object_store(object_store_url.as_ref(), object_store);
228    let reader_factory = lakehouse.reader_factory().clone();
229    register_functions(
230        &ctx,
231        lakehouse.clone(),
232        part_provider.clone(),
233        query_range,
234        view_factory.clone(),
235        is_admin,
236    );
237    for view in view_factory.get_global_views() {
238        register_table(
239            lakehouse.clone(),
240            reader_factory.clone(),
241            part_provider.clone(),
242            query_range,
243            &ctx,
244            view.clone(),
245        )
246        .await?;
247    }
248    // Apply custom configuration
249    configurator.configure(&ctx).await?;
250    Ok(ctx)
251}
252
253#[span_fn]
254pub async fn query(
255    lakehouse: Arc<LakehouseContext>,
256    part_provider: Arc<dyn QueryPartitionProvider>,
257    query_range: Option<TimeRange>,
258    sql: &str,
259    view_factory: Arc<ViewFactory>,
260    configurator: Arc<dyn SessionConfigurator>,
261    is_admin: bool,
262) -> Result<Answer> {
263    info!("query sql={sql}");
264    let ctx = make_session_context(
265        lakehouse,
266        part_provider,
267        query_range,
268        view_factory,
269        configurator,
270        is_admin,
271    )
272    .await
273    .with_context(|| "make_session_context")?;
274    let df = ctx.sql(sql).await?;
275    let schema = df.schema().inner().clone();
276    let batches: Vec<RecordBatch> = df.collect().await?;
277    Ok(Answer::new(schema, batches))
278}