Skip to main content

micromegas_analytics/lakehouse/
block_partition_spec.rs

1use super::{
2    partition_source_data::{PartitionBlocksSource, PartitionSourceBlock},
3    view::{PartitionSpec, ViewMetadata},
4    write_partition::{PartitionRowSet, write_partition_from_rows},
5};
6use crate::{response_writer::Logger, time::TimeRange};
7use anyhow::{Context, Result};
8use async_trait::async_trait;
9use datafusion::arrow::datatypes::Schema;
10use futures::StreamExt;
11use micromegas_ingestion::data_lake_connection::DataLakeConnection;
12use micromegas_telemetry::blob_storage::BlobStorage;
13use micromegas_tracing::prelude::*;
14use std::collections::HashMap;
15use std::fmt::Debug;
16use std::sync::Arc;
17
18/// BlockProcessor transforms a single block of telemetry into a set of rows
19#[async_trait]
20pub trait BlockProcessor: Send + Sync + Debug {
21    /// Processes a single block of telemetry.
22    async fn process(
23        &self,
24        blob_storage: Arc<BlobStorage>,
25        src_block: Arc<PartitionSourceBlock>,
26    ) -> Result<Option<PartitionRowSet>>;
27}
28
29/// Map from `streams.format` to the processor that handles that wire format.
30/// Views register one entry per format they understand (e.g. log entries register
31/// both `"micromegas-transit"` and `"otlp/v1/logs"`).
32pub type BlockProcessorMap = HashMap<&'static str, Arc<dyn BlockProcessor>>;
33
34/// BlockPartitionSpec processes blocks individually and out of order
35/// which works fine for measures & log entries.
36///
37/// Per-block dispatch keys on `PartitionSourceBlock::format` so a single view can
38/// materialize blocks coming from heterogeneous wire formats (native CBOR + OTLP).
39/// Unknown formats are warned and skipped.
40#[derive(Debug)]
41pub struct BlockPartitionSpec {
42    pub view_metadata: ViewMetadata,
43    pub schema: Arc<Schema>,
44    pub insert_range: TimeRange,
45    pub source_data: Arc<dyn PartitionBlocksSource>,
46    pub block_processors: Arc<BlockProcessorMap>,
47}
48
49#[async_trait]
50impl PartitionSpec for BlockPartitionSpec {
51    fn is_empty(&self) -> bool {
52        self.source_data.is_empty()
53    }
54
55    fn get_source_data_hash(&self) -> Vec<u8> {
56        self.source_data.get_source_data_hash()
57    }
58
59    #[span_fn]
60    async fn write(&self, lake: Arc<DataLakeConnection>, logger: Arc<dyn Logger>) -> Result<()> {
61        let desc = format!(
62            "[{}, {}] {} {}",
63            self.view_metadata.view_set_name,
64            self.view_metadata.view_instance_id,
65            self.insert_range.begin.to_rfc3339(),
66            self.insert_range.end.to_rfc3339()
67        );
68        instrument_named!(
69            logger.write_log_entry(format!("writing {desc}")),
70            "log_write_start"
71        )
72        .await?;
73
74        instrument_named!(
75            logger.write_log_entry(format!(
76                "reading {} blocks",
77                self.source_data.get_nb_blocks()
78            )),
79            "log_reading_blocks"
80        )
81        .await?;
82
83        // Allow empty source data - write_partition_from_rows will create
84        // an empty partition record if no data is sent through the channel
85        let (tx, rx) = tokio::sync::mpsc::channel(1);
86        let join_handle = spawn_with_context(write_partition_from_rows(
87            lake.clone(),
88            self.view_metadata.clone(),
89            self.schema.clone(),
90            self.insert_range,
91            self.source_data.get_source_data_hash(),
92            None,
93            rx,
94            logger.clone(),
95        ));
96
97        // If source data is empty, just close the channel to create an empty partition
98        if self.source_data.is_empty() {
99            drop(tx);
100            instrument_named!(join_handle, "write_partition_from_rows_join").await??;
101            return Ok(());
102        }
103
104        let max_size = self.source_data.get_max_payload_size() as usize;
105        let mut nb_tasks = (100 * 1024 * 1024) / max_size; // try to download up to 100 MB of payloads
106        nb_tasks = nb_tasks.clamp(1, 64);
107
108        let mut stream =
109            instrument_named!(self.source_data.get_blocks_stream(), "get_blocks_stream")
110                .await
111                .map(|src_block_res| async {
112                    let src_block = src_block_res.with_context(|| "get_blocks_stream")?;
113                    // Per-block dispatch on `streams.format`. A view that doesn't register
114                    // a processor for some format silently skips matching blocks instead of
115                    // erroring — keeps the partition build moving when an unknown format
116                    // shows up alongside known ones.
117                    let Some(block_processor) = self
118                        .block_processors
119                        .get(src_block.format.as_str())
120                        .cloned()
121                    else {
122                        warn!(
123                            "no block processor for format={} (view={}/{}); skipping block_id={}",
124                            src_block.format,
125                            self.view_metadata.view_set_name,
126                            self.view_metadata.view_instance_id,
127                            src_block.block.block_id
128                        );
129                        return Ok::<Option<PartitionRowSet>, anyhow::Error>(None);
130                    };
131                    let blob_storage = lake.blob_storage.clone();
132                    let handle = spawn_with_context(async move {
133                        block_processor
134                            .process(blob_storage, src_block)
135                            .await
136                            .with_context(|| "processing source block")
137                    });
138                    instrument_named!(handle, "block_processor_task_join")
139                        .await
140                        .with_context(|| "handle.await")?
141                })
142                .buffer_unordered(nb_tasks);
143
144        while let Some(res_opt_rows) = instrument_named!(stream.next(), "blocks_stream_next").await
145        {
146            match res_opt_rows {
147                Err(e) => {
148                    error!("{e:?}");
149                    instrument_named!(logger.write_log_entry(format!("{e:?}")), "log_block_error")
150                        .await?;
151                }
152                Ok(Some(row_set)) => {
153                    instrument_named!(tx.send(Ok(row_set)), "row_set_channel_send").await?;
154                }
155                Ok(None) => {
156                    debug!("empty block");
157                }
158            }
159        }
160        drop(tx);
161        instrument_named!(join_handle, "write_partition_from_rows_join").await??;
162        Ok(())
163    }
164}