micromegas_analytics/lakehouse/
image_block_processor.rs1use 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}