micromegas_analytics/lakehouse/
block_partition_spec.rs1use 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#[async_trait]
20pub trait BlockProcessor: Send + Sync + Debug {
21 async fn process(
23 &self,
24 blob_storage: Arc<BlobStorage>,
25 src_block: Arc<PartitionSourceBlock>,
26 ) -> Result<Option<PartitionRowSet>>;
27}
28
29pub type BlockProcessorMap = HashMap<&'static str, Arc<dyn BlockProcessor>>;
33
34#[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 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 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; 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 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}