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#[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#[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#[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#[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 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 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}