Skip to main content

micromegas_analytics/
replication.rs

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        // Hard failure (rather than a silent default) on missing `format` so a
39        // v3 source replicating into a v4 target surfaces the schema mismatch loudly.
40        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
230/// Ingests data from a FlightRecordBatchStream into the specified table.
231pub 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}