micromegas_analytics/lakehouse/
parse_block_table_function.rs1use super::{
2 lakehouse_context::LakehouseContext, partition_cache::QueryPartitionProvider,
3 session_configurator::NoOpSessionConfigurator, view_factory::ViewFactory,
4};
5use crate::{
6 dfext::{string_column_accessor::string_column_by_name, typed_column::typed_column_by_name},
7 metadata::StreamMetadata,
8 payload::{fetch_block_payload, parse_block},
9 time::TimeRange,
10};
11use anyhow::Context;
12use async_trait::async_trait;
13use datafusion::{
14 arrow::{
15 array::{BinaryBuilder, Int64Array, Int64Builder, RecordBatch, StringBuilder},
16 datatypes::{DataType, Field, Schema, SchemaRef},
17 },
18 catalog::{Session, TableFunctionArgs, TableFunctionImpl, TableProvider},
19 common::plan_err,
20 datasource::{
21 TableType,
22 memory::{DataSourceExec, MemorySourceConfig},
23 },
24 error::DataFusionError,
25 physical_plan::ExecutionPlan,
26 prelude::Expr,
27};
28use jsonb::Value as JsonbValue;
29use micromegas_ingestion::web_ingestion_service::FORMAT_TRANSIT;
30use micromegas_tracing::prelude::*;
31use micromegas_transit::{UserDefinedType, value::Value as TransitValue};
32use std::{borrow::Cow, collections::BTreeMap, sync::Arc};
33use uuid::Uuid;
34
35use crate::dfext::expressions::exp_to_string;
36
37fn output_schema() -> SchemaRef {
38 Arc::new(Schema::new(vec![
39 Field::new("object_index", DataType::Int64, false),
40 Field::new("type_name", DataType::Utf8, false),
41 Field::new("value", DataType::Binary, false),
42 ]))
43}
44
45pub fn transit_value_to_jsonb(value: TransitValue<'_>) -> JsonbValue<'_> {
47 match value {
48 TransitValue::String(s) => JsonbValue::String(Cow::Borrowed(s)),
49 TransitValue::Object(obj) => {
50 let mut map = BTreeMap::new();
51 map.insert(
52 "__type".to_string(),
53 JsonbValue::String(Cow::Borrowed(obj.type_name)),
54 );
55 for &(name, val) in obj.members {
56 map.insert(name.to_string(), transit_value_to_jsonb(val));
57 }
58 JsonbValue::Object(map)
59 }
60 TransitValue::U8(v) => JsonbValue::Number(jsonb::Number::UInt64(u64::from(v))),
61 TransitValue::U32(v) => JsonbValue::Number(jsonb::Number::UInt64(u64::from(v))),
62 TransitValue::U64(v) => JsonbValue::Number(jsonb::Number::UInt64(v)),
63 TransitValue::I64(v) => JsonbValue::Number(jsonb::Number::Int64(v)),
64 TransitValue::F64(v) => JsonbValue::Number(jsonb::Number::Float64(v)),
65 TransitValue::None => JsonbValue::Null,
66 TransitValue::Bytes(b) => {
67 JsonbValue::String(Cow::Owned(format!("<binary {} bytes>", b.len())))
68 }
69 }
70}
71
72async fn fetch_block_metadata(
75 lakehouse: Arc<LakehouseContext>,
76 part_provider: Arc<dyn QueryPartitionProvider>,
77 query_range: Option<TimeRange>,
78 view_factory: Arc<ViewFactory>,
79 block_id_str: &str,
80) -> anyhow::Result<Option<(Uuid, i64, StreamMetadata)>> {
81 let ctx = super::query::make_session_context(
82 lakehouse,
83 part_provider,
84 query_range,
85 view_factory,
86 Arc::new(NoOpSessionConfigurator),
87 false,
88 )
89 .await?;
90
91 let sql = format!(
92 "SELECT block_id, stream_id, process_id, object_offset,
93 \"streams.dependencies_metadata\", \"streams.objects_metadata\", \"streams.format\"
94 FROM blocks
95 WHERE block_id = '{block_id_str}'"
96 );
97 let df = ctx.sql(&sql).await?;
98 let batches = df.collect().await?;
99
100 if batches.is_empty() || batches[0].num_rows() == 0 {
101 return Ok(None);
102 }
103
104 let batch = &batches[0];
105
106 let block_id_col = string_column_by_name(batch, "block_id")?;
107 let stream_id_col = string_column_by_name(batch, "stream_id")?;
108 let process_id_col = string_column_by_name(batch, "process_id")?;
109 let object_offset_col: &Int64Array = typed_column_by_name(batch, "object_offset")?;
110 let format_col = string_column_by_name(batch, "streams.format")?;
111 let format = format_col.value(0)?;
112 if format != FORMAT_TRANSIT {
113 anyhow::bail!(
114 "parse_block does not support format={format} (only {FORMAT_TRANSIT}). \
115 Query `log_entries`/`measures`/`otel_spans` instead for OTel data."
116 );
117 }
118
119 let block_id = Uuid::parse_str(block_id_col.value(0)?)?;
120 let stream_id = Uuid::parse_str(stream_id_col.value(0)?)?;
121 let process_id = Uuid::parse_str(process_id_col.value(0)?)?;
122 let object_offset = object_offset_col.value(0);
123
124 let deps_col = batch
125 .column_by_name("streams.dependencies_metadata")
126 .context("streams.dependencies_metadata column not found")?;
127 let deps_binary: &datafusion::arrow::array::BinaryArray = deps_col
128 .as_any()
129 .downcast_ref()
130 .context("failed to cast dependencies_metadata to BinaryArray")?;
131 let deps_bytes = deps_binary.value(0);
132 let dependencies_metadata: Vec<UserDefinedType> =
133 ciborium::from_reader(deps_bytes).context("decoding dependencies_metadata")?;
134
135 let objs_col = batch
136 .column_by_name("streams.objects_metadata")
137 .context("streams.objects_metadata column not found")?;
138 let objs_binary: &datafusion::arrow::array::BinaryArray = objs_col
139 .as_any()
140 .downcast_ref()
141 .context("failed to cast objects_metadata to BinaryArray")?;
142 let objs_bytes = objs_binary.value(0);
143 let objects_metadata: Vec<UserDefinedType> =
144 ciborium::from_reader(objs_bytes).context("decoding objects_metadata")?;
145
146 let stream_metadata = StreamMetadata {
147 process_id,
148 stream_id,
149 dependencies_metadata,
150 objects_metadata,
151 tags: vec![],
152 properties: Arc::new(vec![]),
153 };
154
155 Ok(Some((block_id, object_offset, stream_metadata)))
156}
157
158fn parse_block_objects(
160 stream_metadata: &StreamMetadata,
161 payload: µmegas_telemetry::block_wire_format::BlockPayload,
162 object_offset: i64,
163 early_limit: Option<usize>,
164) -> anyhow::Result<RecordBatch> {
165 let mut index_builder = Int64Builder::new();
166 let mut name_builder = StringBuilder::new();
167 let mut value_builder = BinaryBuilder::new();
168 let mut local_index: i64 = 0;
169 let mut nb_objects: usize = 0;
170
171 parse_block(stream_metadata, payload, |value| {
172 if let TransitValue::Object(obj) = value {
173 let jsonb_val = transit_value_to_jsonb(value);
174 let mut buf = Vec::new();
175 jsonb_val.write_to_vec(&mut buf);
176
177 index_builder.append_value(object_offset + local_index);
178 name_builder.append_value(obj.type_name);
179 value_builder.append_value(&buf);
180 nb_objects += 1;
181 } else {
182 warn!(
183 "parse_block: skipping non-Object value at index {}",
184 object_offset + local_index
185 );
186 }
187 local_index += 1;
188
189 if let Some(lim) = early_limit {
190 Ok(nb_objects < lim)
191 } else {
192 Ok(true)
193 }
194 })?;
195
196 Ok(RecordBatch::try_new(
197 output_schema(),
198 vec![
199 Arc::new(index_builder.finish()),
200 Arc::new(name_builder.finish()),
201 Arc::new(value_builder.finish()),
202 ],
203 )?)
204}
205
206#[derive(Debug)]
209pub struct ParseBlockTableFunction {
210 lakehouse: Arc<LakehouseContext>,
211 view_factory: Arc<ViewFactory>,
212 part_provider: Arc<dyn QueryPartitionProvider>,
213 query_range: Option<TimeRange>,
214}
215
216impl ParseBlockTableFunction {
217 pub fn new(
218 lakehouse: Arc<LakehouseContext>,
219 view_factory: Arc<ViewFactory>,
220 part_provider: Arc<dyn QueryPartitionProvider>,
221 query_range: Option<TimeRange>,
222 ) -> Self {
223 Self {
224 lakehouse,
225 view_factory,
226 part_provider,
227 query_range,
228 }
229 }
230}
231
232impl TableFunctionImpl for ParseBlockTableFunction {
233 fn call_with_args(
234 &self,
235 args: TableFunctionArgs,
236 ) -> datafusion::error::Result<Arc<dyn TableProvider>> {
237 let exprs = args.exprs();
238 let arg = exprs.first().map(exp_to_string);
239 let Some(Ok(block_id)) = arg else {
240 return plan_err!(
241 "First argument to parse_block must be a string (the block ID), given {:?}",
242 arg
243 );
244 };
245 Ok(Arc::new(ParseBlockProvider {
246 block_id,
247 lakehouse: self.lakehouse.clone(),
248 view_factory: self.view_factory.clone(),
249 part_provider: self.part_provider.clone(),
250 query_range: self.query_range,
251 }))
252 }
253}
254
255#[derive(Debug)]
256struct ParseBlockProvider {
257 block_id: String,
258 lakehouse: Arc<LakehouseContext>,
259 view_factory: Arc<ViewFactory>,
260 part_provider: Arc<dyn QueryPartitionProvider>,
261 query_range: Option<TimeRange>,
262}
263
264#[async_trait]
265impl TableProvider for ParseBlockProvider {
266 fn schema(&self) -> SchemaRef {
267 output_schema()
268 }
269
270 fn table_type(&self) -> TableType {
271 TableType::Temporary
272 }
273
274 async fn scan(
275 &self,
276 _state: &dyn Session,
277 projection: Option<&Vec<usize>>,
278 filters: &[Expr],
279 limit: Option<usize>,
280 ) -> datafusion::error::Result<Arc<dyn ExecutionPlan>> {
281 let block_id_str = &self.block_id;
282
283 let Some((block_id, object_offset, stream_metadata)) = fetch_block_metadata(
284 self.lakehouse.clone(),
285 self.part_provider.clone(),
286 self.query_range,
287 self.view_factory.clone(),
288 block_id_str,
289 )
290 .await
291 .map_err(|e| DataFusionError::External(e.into()))?
292 else {
293 let source = MemorySourceConfig::try_new(
294 &[vec![]],
295 self.schema(),
296 projection.map(|v| v.to_owned()),
297 )?;
298 return Ok(DataSourceExec::from_data_source(source));
299 };
300
301 let blob_storage = self.lakehouse.lake().blob_storage.clone();
303 let payload = fetch_block_payload(
304 blob_storage,
305 sqlx::types::Uuid::from_bytes(*stream_metadata.process_id.as_bytes()),
306 sqlx::types::Uuid::from_bytes(*stream_metadata.stream_id.as_bytes()),
307 sqlx::types::Uuid::from_bytes(*block_id.as_bytes()),
308 )
309 .await
310 .map_err(|e| DataFusionError::External(e.into()))?;
311
312 let early_limit = if filters.is_empty() { limit } else { None };
314 let rb = parse_block_objects(&stream_metadata, &payload, object_offset, early_limit)
315 .with_context(|| format!("parsing block {block_id_str}"))
316 .map_err(|e| DataFusionError::External(e.into()))?;
317
318 let source = MemorySourceConfig::try_new(
319 &[vec![rb]],
320 self.schema(),
321 projection.map(|v| v.to_owned()),
322 )?;
323 Ok(DataSourceExec::from_data_source(source))
324 }
325}