1use anyhow::{Context, Result};
2use arrow_flight::decode::FlightRecordBatchStream;
3use chrono::DateTime;
4use datafusion::arrow::array::{
5 BinaryArray, GenericListArray, Int32Array, Int64Array, StringArray, TimestampNanosecondArray,
6};
7use futures::StreamExt;
8use micromegas_ingestion::data_lake_connection::DataLakeConnection;
9use micromegas_tracing::prelude::*;
10use std::sync::Arc;
11use uuid::Uuid;
12
13use crate::{
14 dfext::{string_column_accessor::string_column_by_name, typed_column::typed_column_by_name},
15 properties::{
16 properties_column_accessor::properties_column_by_name,
17 utils::extract_properties_from_properties_column,
18 },
19};
20async fn ingest_streams(
21 lake: Arc<DataLakeConnection>,
22 mut rb_stream: FlightRecordBatchStream,
23) -> Result<i64> {
24 let mut tr = lake.db_pool.begin().await?;
25 let mut nb_rows: i64 = 0;
26 while let Some(res) = rb_stream.next().await {
27 let b = res?;
28 nb_rows += b.num_rows() as i64;
29 let stream_id_column = string_column_by_name(&b, "stream_id")?;
30 let process_id_column = string_column_by_name(&b, "process_id")?;
31 let dependencies_metadata_column: &BinaryArray =
32 typed_column_by_name(&b, "dependencies_metadata")?;
33 let objects_metadata_column: &BinaryArray = typed_column_by_name(&b, "objects_metadata")?;
34 let tags_column: &GenericListArray<i32> = typed_column_by_name(&b, "tags")?;
35 let properties_accessor = properties_column_by_name(&b, "properties")?;
36 let insert_time_column: &TimestampNanosecondArray =
37 typed_column_by_name(&b, "insert_time")?;
38 let format_column = string_column_by_name(&b, "format")?;
41
42 for row in 0..b.num_rows() {
43 let stream_id = Uuid::parse_str(stream_id_column.value(row)?)?;
44 let process_id = Uuid::parse_str(process_id_column.value(row)?)?;
45 let tags: Vec<String> = tags_column
46 .value(row)
47 .as_any()
48 .downcast_ref::<StringArray>()
49 .with_context(|| "casting tags")?
50 .iter()
51 .map(|item| String::from(item.unwrap_or_default()))
52 .collect();
53 let properties_map =
54 extract_properties_from_properties_column(properties_accessor.as_ref(), row)?;
55 let properties = micromegas_telemetry::property::make_properties(&properties_map);
56
57 instrument_named!(
58 sqlx::query(
59 "INSERT INTO streams (stream_id, process_id, dependencies_metadata, objects_metadata, tags, properties, insert_time, format)
60 VALUES ($1,$2,$3,$4,$5,$6,$7,$8)
61 ON CONFLICT (stream_id) DO NOTHING;",
62 )
63 .bind(stream_id)
64 .bind(process_id)
65 .bind(dependencies_metadata_column.value(row))
66 .bind(objects_metadata_column.value(row))
67 .bind(tags)
68 .bind(properties)
69 .bind(DateTime::from_timestamp_nanos(
70 insert_time_column.value(row),
71 ))
72 .bind(format_column.value(row)?)
73 .execute(&mut *tr),
74 "sql_insert_stream"
75 )
76 .await
77 .with_context(|| "inserting into streams")?;
78 }
79 }
80 tr.commit().await?;
81 info!("ingested {nb_rows} streams");
82 Ok(nb_rows)
83}
84
85async fn ingest_processes(
86 lake: Arc<DataLakeConnection>,
87 mut rb_stream: FlightRecordBatchStream,
88) -> Result<i64> {
89 let mut tr = lake.db_pool.begin().await?;
90 let mut nb_rows: i64 = 0;
91 while let Some(res) = rb_stream.next().await {
92 let b = res?;
93 nb_rows += b.num_rows() as i64;
94 let process_id_column = string_column_by_name(&b, "process_id")?;
95 let exe_column = string_column_by_name(&b, "exe")?;
96 let username_column = string_column_by_name(&b, "username")?;
97 let realname_column = string_column_by_name(&b, "realname")?;
98 let computer_column = string_column_by_name(&b, "computer")?;
99 let distro_column = string_column_by_name(&b, "distro")?;
100 let cpu_brand_column = string_column_by_name(&b, "cpu_brand")?;
101 let process_tsc_freq_column: &Int64Array = typed_column_by_name(&b, "tsc_frequency")?;
102 let start_time_column: &TimestampNanosecondArray = typed_column_by_name(&b, "start_time")?;
103 let start_ticks_column: &Int64Array = typed_column_by_name(&b, "start_ticks")?;
104 let insert_time_column: &TimestampNanosecondArray =
105 typed_column_by_name(&b, "insert_time")?;
106 let parent_process_id_column = string_column_by_name(&b, "parent_process_id")?;
107 let properties_accessor = properties_column_by_name(&b, "properties")?;
108 for row in 0..b.num_rows() {
109 let process_id = Uuid::parse_str(process_id_column.value(row)?)?;
110 let parent_process_str = parent_process_id_column.value(row)?;
111 let parent_process_id = if parent_process_str.is_empty() {
112 None
113 } else {
114 Some(Uuid::parse_str(parent_process_str)?)
115 };
116 let properties_map =
117 extract_properties_from_properties_column(properties_accessor.as_ref(), row)?;
118 let properties = micromegas_telemetry::property::make_properties(&properties_map);
119 instrument_named!(
120 sqlx::query(
121 "INSERT INTO processes VALUES($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13) ON CONFLICT (process_id) DO NOTHING;",
122 )
123 .bind(process_id)
124 .bind(exe_column.value(row)?)
125 .bind(username_column.value(row)?)
126 .bind(realname_column.value(row)?)
127 .bind(computer_column.value(row)?)
128 .bind(distro_column.value(row)?)
129 .bind(cpu_brand_column.value(row)?)
130 .bind(process_tsc_freq_column.value(row))
131 .bind(DateTime::from_timestamp_nanos(start_time_column.value(row)))
132 .bind(start_ticks_column.value(row))
133 .bind(DateTime::from_timestamp_nanos(
134 insert_time_column.value(row),
135 ))
136 .bind(parent_process_id)
137 .bind(properties)
138 .execute(&mut *tr),
139 "sql_insert_process"
140 )
141 .await
142 .with_context(|| "executing sql insert into processes")?;
143 }
144 }
145 tr.commit().await?;
146 info!("ingested {nb_rows} processes");
147 Ok(nb_rows)
148}
149
150async fn ingest_payloads(
151 lake: Arc<DataLakeConnection>,
152 mut rb_stream: FlightRecordBatchStream,
153) -> Result<i64> {
154 let mut nb_rows: i64 = 0;
155 while let Some(res) = rb_stream.next().await {
156 let b = res?;
157 nb_rows += b.num_rows() as i64;
158 let process_id_column = string_column_by_name(&b, "process_id")?;
159 let stream_id_column = string_column_by_name(&b, "stream_id")?;
160 let block_id_column = string_column_by_name(&b, "block_id")?;
161 let payload_column: &BinaryArray = typed_column_by_name(&b, "payload")?;
162 for row in 0..b.num_rows() {
163 let process_id = process_id_column.value(row)?;
164 let stream_id = stream_id_column.value(row)?;
165 let block_id = block_id_column.value(row)?;
166 let obj_path = format!("blobs/{process_id}/{stream_id}/{block_id}");
167 let payload = bytes::Bytes::copy_from_slice(payload_column.value(row));
168 lake.blob_storage
169 .put(&obj_path, payload)
170 .await
171 .with_context(|| "Error writing block to blob storage")?;
172 }
173 }
174 info!("ingested {nb_rows} payloads");
175 Ok(nb_rows)
176}
177
178async fn ingest_blocks(
179 lake: Arc<DataLakeConnection>,
180 mut rb_stream: FlightRecordBatchStream,
181) -> Result<i64> {
182 let mut tr = lake.db_pool.begin().await?;
183 let mut nb_rows: i64 = 0;
184 while let Some(res) = rb_stream.next().await {
185 let b = res?;
186 nb_rows += b.num_rows() as i64;
187 let block_id_column = string_column_by_name(&b, "block_id")?;
188 let stream_id_column = string_column_by_name(&b, "stream_id")?;
189 let process_id_column = string_column_by_name(&b, "process_id")?;
190 let begin_time_column: &TimestampNanosecondArray = typed_column_by_name(&b, "begin_time")?;
191 let begin_ticks_column: &Int64Array = typed_column_by_name(&b, "begin_ticks")?;
192 let end_time_column: &TimestampNanosecondArray = typed_column_by_name(&b, "end_time")?;
193 let end_ticks_column: &Int64Array = typed_column_by_name(&b, "end_ticks")?;
194 let nb_objects_column: &Int32Array = typed_column_by_name(&b, "nb_objects")?;
195 let object_offset_column: &Int64Array = typed_column_by_name(&b, "object_offset")?;
196 let payload_size_column: &Int64Array = typed_column_by_name(&b, "payload_size")?;
197 let insert_time_column: &TimestampNanosecondArray =
198 typed_column_by_name(&b, "insert_time")?;
199 for row in 0..b.num_rows() {
200 let block_id = Uuid::parse_str(block_id_column.value(row)?)?;
201 let stream_id = Uuid::parse_str(stream_id_column.value(row)?)?;
202 let process_id = Uuid::parse_str(process_id_column.value(row)?)?;
203 instrument_named!(
204 sqlx::query("INSERT INTO blocks VALUES($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11) ON CONFLICT (block_id) DO NOTHING;")
205 .bind(block_id)
206 .bind(stream_id)
207 .bind(process_id)
208 .bind(DateTime::from_timestamp_nanos(begin_time_column.value(row)))
209 .bind(begin_ticks_column.value(row))
210 .bind(DateTime::from_timestamp_nanos(end_time_column.value(row)))
211 .bind(end_ticks_column.value(row))
212 .bind(nb_objects_column.value(row))
213 .bind(object_offset_column.value(row))
214 .bind(payload_size_column.value(row))
215 .bind(DateTime::from_timestamp_nanos(
216 insert_time_column.value(row),
217 ))
218 .execute(&mut *tr),
219 "sql_insert_block"
220 )
221 .await
222 .with_context(|| "executing sql insert into blocks")?;
223 }
224 }
225 tr.commit().await?;
226 info!("ingested {nb_rows} blocks");
227 Ok(nb_rows)
228}
229
230pub async fn bulk_ingest(
232 lake: Arc<DataLakeConnection>,
233 table_name: &str,
234 rb_stream: FlightRecordBatchStream,
235) -> Result<i64> {
236 match table_name {
237 "processes" => ingest_processes(lake, rb_stream).await,
238 "streams" => ingest_streams(lake, rb_stream).await,
239 "blocks" => ingest_blocks(lake, rb_stream).await,
240 "payloads" => ingest_payloads(lake, rb_stream).await,
241 other => anyhow::bail!("bulk ingest for table {other} not supported"),
242 }
243}