Skip to main content

micromegas_analytics/lakehouse/
image_block_processor.rs

1use super::{
2    block_partition_spec::BlockProcessor, partition_source_data::PartitionSourceBlock,
3    write_partition::PartitionRowSet,
4};
5use crate::{
6    images_table::ImagesRecordBuilder,
7    payload::{fetch_block_payload, parse_block},
8    time::make_time_converter_from_block_meta,
9};
10use anyhow::{Context, Result};
11use async_trait::async_trait;
12use micromegas_telemetry::blob_storage::BlobStorage;
13use micromegas_tracing::prelude::*;
14use micromegas_transit::value::Value;
15use std::sync::Arc;
16
17#[derive(Debug)]
18pub struct ImageBlockProcessor {}
19
20#[async_trait]
21impl BlockProcessor for ImageBlockProcessor {
22    #[span_fn]
23    async fn process(
24        &self,
25        blob_storage: Arc<BlobStorage>,
26        src_block: Arc<PartitionSourceBlock>,
27    ) -> Result<Option<PartitionRowSet>> {
28        let convert_ticks =
29            make_time_converter_from_block_meta(&src_block.process, &src_block.block)?;
30        let payload = fetch_block_payload(
31            blob_storage,
32            src_block.process.process_id,
33            src_block.stream.stream_id,
34            src_block.block.block_id,
35        )
36        .await
37        .with_context(|| "fetch_block_payload")?;
38
39        let process_id_str = format!("{}", src_block.process.process_id);
40        let stream_id_str = format!("{}", src_block.stream.stream_id);
41        let block_id_str = format!("{}", src_block.block.block_id);
42        let insert_time_nanos = src_block
43            .block
44            .insert_time
45            .timestamp_nanos_opt()
46            .with_context(|| "converting insert_time to nanoseconds")?;
47
48        let mut record_builder = ImagesRecordBuilder::new();
49
50        parse_block(&src_block.stream, &payload, |val| {
51            if let Value::Object(obj) = val
52                && obj.type_name == "ImageEvent"
53            {
54                let ticks = obj.get::<i64>("time").with_context(|| "reading time")?;
55                let name = obj.get::<&str>("name").with_context(|| "reading name")?;
56                let format = obj
57                    .get::<&str>("format")
58                    .with_context(|| "reading format")?;
59                let image_data = obj.get::<&[u8]>("data").with_context(|| "reading data")?;
60                let time_ns = convert_ticks.ticks_to_nanoseconds(ticks);
61                let payload_size = image_data.len() as i64;
62                record_builder.append(
63                    &src_block.process,
64                    &process_id_str,
65                    &stream_id_str,
66                    &block_id_str,
67                    insert_time_nanos,
68                    time_ns,
69                    name,
70                    format,
71                    payload_size,
72                    image_data,
73                )?;
74            }
75            Ok(true)
76        })
77        .with_context(|| "parse_block")?;
78
79        if let Some(time_range) = record_builder.get_time_range() {
80            let record_batch = record_builder.finish()?;
81            Ok(Some(PartitionRowSet {
82                rows_time_range: time_range,
83                rows: record_batch,
84            }))
85        } else {
86            Ok(None)
87        }
88    }
89}