Skip to main content

micromegas_analytics/lakehouse/
parse_block_table_function.rs

1use 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
45/// Converts a `transit::Value` to a `jsonb::Value`.
46pub 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
72/// Queries the global blocks view for a block's metadata and constructs a `StreamMetadata`.
73/// Returns `None` if the block is not found.
74async 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
158/// Parses transit objects from a block payload and returns them as a RecordBatch.
159fn parse_block_objects(
160    stream_metadata: &StreamMetadata,
161    payload: &micromegas_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/// A DataFusion `TableFunctionImpl` that parses a block's transit-serialized objects
207/// and returns each object as a row with its type name and full content as JSONB.
208#[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        // Fetch and parse the block payload
302        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        // Parse transit objects and convert to JSONB
313        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}