1use super::{
2 block_partition_spec::{BlockPartitionSpec, BlockProcessorMap},
3 blocks_view::BlocksView,
4 lakehouse_context::LakehouseContext,
5 partition_cache::{LivePartitionProvider, QueryPartitionProvider},
6 partition_source_data::{PartitionSourceBlock, SourceDataBlocksInMemory},
7 view::{View, ViewMetadata},
8};
9use crate::{
10 dfext::{
11 string_column_accessor::string_column_by_name,
12 typed_column::{get_single_row_primitive_value, typed_column_by_name},
13 },
14 lakehouse::{
15 partition_cache::PartitionCache, partition_source_data::hash_to_object_count,
16 query::query_partitions, view::PartitionSpec,
17 },
18 metadata::{ProcessMetadata, StreamMetadata, block_from_batch_row},
19 properties::properties_column_accessor::properties_column_by_name,
20 response_writer::ResponseWriter,
21 time::TimeRange,
22};
23use anyhow::{Context, Result};
24use chrono::DurationRound;
25use chrono::{DateTime, TimeDelta, Utc};
26use datafusion::arrow::array::{BinaryArray, GenericListArray, StringArray};
27use datafusion::arrow::datatypes::{Schema, TimestampNanosecondType};
28use micromegas_ingestion::data_lake_connection::DataLakeConnection;
29use micromegas_tracing::prelude::*;
30use sqlx::Row;
31use std::sync::Arc;
32use uuid::Uuid;
33
34pub struct JitPartitionConfig {
36 pub max_nb_objects: i64,
37 pub max_insert_time_slice: TimeDelta,
38}
39
40impl Default for JitPartitionConfig {
41 fn default() -> Self {
42 JitPartitionConfig {
43 max_nb_objects: 20 * 1024 * 1024,
44 max_insert_time_slice: TimeDelta::hours(1),
45 }
46 }
47}
48
49async fn get_insert_time_range(
50 lakehouse: Arc<LakehouseContext>,
51 blocks_view: &BlocksView,
52 query_time_range: &TimeRange,
53 stream: Arc<StreamMetadata>,
54) -> Result<Option<TimeRange>> {
55 let part_provider = LivePartitionProvider::new(lakehouse.lake().db_pool.clone());
57 let partitions = part_provider
58 .fetch(
59 &blocks_view.get_view_set_name(),
60 &blocks_view.get_view_instance_id(),
61 Some(*query_time_range),
62 blocks_view.get_file_schema_hash(),
63 )
64 .await?;
65 let stream_id = &stream.stream_id;
66 let begin_range_iso = query_time_range.begin.to_rfc3339();
67 let end_range_iso = query_time_range.end.to_rfc3339();
68 let sql = format!(
69 r#"SELECT MIN(insert_time) as min_insert_time, MAX(insert_time) as max_insert_time
70 FROM source
71 WHERE stream_id = '{stream_id}'
72 AND begin_time <= '{end_range_iso}'
73 AND end_time >= '{begin_range_iso}';"#
74 );
75 let reader_factory = lakehouse.reader_factory().clone();
76 let rbs = query_partitions(
77 lakehouse.runtime().clone(),
78 reader_factory,
79 lakehouse.lake().blob_storage.inner(),
80 blocks_view.get_file_schema(),
81 Arc::new(partitions),
82 &sql,
83 )
84 .await?
85 .collect()
86 .await?;
87 if rbs.is_empty() {
88 return Ok(None);
89 }
90 if rbs[0].num_rows() == 0 {
91 return Ok(None);
92 }
93 let min_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 0)?;
94 let max_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 1)?;
95 Ok(Some(TimeRange::new(
96 DateTime::from_timestamp_nanos(min_insert_time),
97 DateTime::from_timestamp_nanos(max_insert_time),
98 )))
99}
100
101pub async fn generate_stream_jit_partitions_segment(
103 config: &JitPartitionConfig,
104 lakehouse: Arc<LakehouseContext>,
105 blocks_view: &BlocksView,
106 partitions: &PartitionCache,
107 insert_time_range: &TimeRange,
108 stream: Arc<StreamMetadata>,
109 process: Arc<ProcessMetadata>,
110) -> Result<Vec<SourceDataBlocksInMemory>> {
111 let partitions = partitions
112 .filter_insert_range(*insert_time_range)
113 .partitions;
114
115 let stream_id = &stream.stream_id;
116 let begin_range_iso = insert_time_range.begin.to_rfc3339();
117 let end_range_iso = insert_time_range.end.to_rfc3339();
118 let sql = format!(
119 r#"SELECT block_id, stream_id, process_id, begin_time, end_time, begin_ticks, end_ticks, nb_objects, object_offset, payload_size, insert_time, "streams.format"
120 FROM source
121 WHERE stream_id = '{stream_id}'
122 AND insert_time >= '{begin_range_iso}'
123 AND insert_time < '{end_range_iso}'
124 ORDER BY insert_time, block_id;"#
125 );
126
127 let reader_factory = lakehouse.reader_factory().clone();
128 let rbs = query_partitions(
129 lakehouse.runtime().clone(),
130 reader_factory,
131 lakehouse.lake().blob_storage.inner(),
132 blocks_view.get_file_schema(),
133 Arc::new(partitions),
134 &sql,
135 )
136 .await?
137 .collect()
138 .await?;
139
140 let mut partitions = vec![];
141 let mut partition_blocks = vec![];
142 let mut partition_nb_objects: i64 = 0;
143 for rb in rbs {
144 let format_column = string_column_by_name(&rb, "streams.format")?;
145 for ir in 0..rb.num_rows() {
146 let block = block_from_batch_row(&rb, ir).with_context(|| "block_from_batch_row")?;
147 let block_nb_objects = block.nb_objects as i64;
148 let format = format_column.value(ir)?.to_string();
149
150 if partition_nb_objects + block_nb_objects > config.max_nb_objects
152 && !partition_blocks.is_empty()
153 {
154 partitions.push(SourceDataBlocksInMemory {
156 blocks: partition_blocks,
157 block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
158 });
159 partition_blocks = vec![Arc::new(PartitionSourceBlock {
161 block,
162 stream: stream.clone(),
163 process: process.clone(),
164 format,
165 })];
166 partition_nb_objects = block_nb_objects;
167 } else {
168 partition_nb_objects += block_nb_objects;
170 partition_blocks.push(Arc::new(PartitionSourceBlock {
171 block,
172 stream: stream.clone(),
173 process: process.clone(),
174 format,
175 }));
176 }
177 }
178 }
179 if partition_nb_objects != 0 {
180 partitions.push(SourceDataBlocksInMemory {
181 blocks: partition_blocks,
182 block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
183 });
184 }
185
186 Ok(partitions)
187}
188
189pub async fn generate_stream_jit_partitions(
193 config: &JitPartitionConfig,
194 lakehouse: Arc<LakehouseContext>,
195 blocks_view: &BlocksView,
196 query_time_range: &TimeRange,
197 stream: Arc<StreamMetadata>,
198 process: Arc<ProcessMetadata>,
199) -> Result<Vec<SourceDataBlocksInMemory>> {
200 let insert_time_range = get_insert_time_range(
201 lakehouse.clone(),
202 blocks_view,
203 query_time_range,
204 stream.clone(),
205 )
206 .await?;
207 if insert_time_range.is_none() {
208 return Ok(vec![]);
209 }
210 let insert_time_range = insert_time_range.with_context(|| "missing insert_time_range")?;
211 let insert_time_range = TimeRange::new(
212 insert_time_range
213 .begin
214 .duration_trunc(config.max_insert_time_slice)?,
215 insert_time_range
216 .end
217 .duration_trunc(config.max_insert_time_slice)?
218 + config.max_insert_time_slice,
219 );
220 let segment_source_partitions = instrument_named!(
221 PartitionCache::fetch_overlapping_insert_range_for_view(
222 &lakehouse.lake().db_pool,
223 blocks_view.get_view_set_name(),
224 blocks_view.get_view_instance_id(),
225 insert_time_range,
226 ),
227 "fetch_overlapping_insert_range_for_view"
228 )
229 .await?;
230
231 let mut begin_segment = insert_time_range.begin;
232 let mut end_segment = begin_segment + config.max_insert_time_slice;
233 let mut partitions = vec![];
234 while end_segment <= insert_time_range.end {
235 let insert_time_range = TimeRange::new(begin_segment, end_segment);
236 let mut segment_partitions = generate_stream_jit_partitions_segment(
237 config,
238 lakehouse.clone(),
239 blocks_view,
240 &segment_source_partitions,
241 &insert_time_range,
242 stream.clone(),
243 process.clone(),
244 )
245 .await?;
246 partitions.append(&mut segment_partitions);
247 begin_segment = end_segment;
248 end_segment = begin_segment + config.max_insert_time_slice;
249 }
250 Ok(partitions)
251}
252
253#[span_fn]
255pub async fn generate_process_jit_partitions_segment(
256 config: &JitPartitionConfig,
257 lakehouse: Arc<LakehouseContext>,
258 blocks_view: &BlocksView,
259 partitions: &PartitionCache,
260 insert_time_range: &TimeRange,
261 process: Arc<ProcessMetadata>,
262 stream_tag: &str,
263) -> Result<Vec<SourceDataBlocksInMemory>> {
264 let partitions = partitions
265 .filter_insert_range(*insert_time_range)
266 .partitions;
267
268 let process_id = &process.process_id;
269 let begin_range_iso = insert_time_range.begin.to_rfc3339();
270 let end_range_iso = insert_time_range.end.to_rfc3339();
271 let sql = format!(
272 r#"SELECT block_id, stream_id, process_id, begin_time, end_time, begin_ticks, end_ticks, nb_objects, object_offset, payload_size, insert_time,
273 "streams.dependencies_metadata", "streams.objects_metadata", "streams.tags", "streams.properties", "streams.format"
274 FROM source
275 WHERE process_id = '{process_id}'
276 AND array_has( "streams.tags", '{stream_tag}' )
277 AND insert_time >= '{begin_range_iso}'
278 AND insert_time < '{end_range_iso}'
279 ORDER BY insert_time, block_id;"#
280 );
281
282 let reader_factory = lakehouse.reader_factory().clone();
283 let df = instrument_named!(
284 query_partitions(
285 lakehouse.runtime().clone(),
286 reader_factory,
287 lakehouse.lake().blob_storage.inner(),
288 blocks_view.get_file_schema(),
289 Arc::new(partitions),
290 &sql,
291 ),
292 "query_partitions"
293 )
294 .await?;
295 let rbs = instrument_named!(df.collect(), "collect_partition_blocks").await?;
296
297 let mut partitions = vec![];
298 let mut partition_blocks = vec![];
299 let mut partition_nb_objects: i64 = 0;
300
301 for rb in rbs {
302 for ir in 0..rb.num_rows() {
303 let block = block_from_batch_row(&rb, ir).with_context(|| "block_from_batch_row")?;
304 let block_nb_objects = block.nb_objects as i64;
305
306 let stream_id_column = string_column_by_name(&rb, "stream_id")?;
308 let stream_process_id_column = string_column_by_name(&rb, "process_id")?;
309 let dependencies_metadata_column: &BinaryArray =
310 typed_column_by_name(&rb, "streams.dependencies_metadata")?;
311 let objects_metadata_column: &BinaryArray =
312 typed_column_by_name(&rb, "streams.objects_metadata")?;
313 let stream_tags_column: &GenericListArray<i32> =
314 typed_column_by_name(&rb, "streams.tags")?;
315 let stream_properties_accessor = properties_column_by_name(&rb, "streams.properties")?;
316 let stream_format_column = string_column_by_name(&rb, "streams.format")?;
317
318 let stream_id = Uuid::parse_str(stream_id_column.value(ir)?)
319 .with_context(|| "parsing stream_id")?;
320 let stream_process_id = Uuid::parse_str(stream_process_id_column.value(ir)?)
321 .with_context(|| "parsing stream process_id")?;
322
323 let dependencies_metadata = dependencies_metadata_column.value(ir);
324 let objects_metadata = objects_metadata_column.value(ir);
325 let stream_tags = stream_tags_column
326 .value(ir)
327 .as_any()
328 .downcast_ref::<StringArray>()
329 .with_context(|| "casting stream_tags")?
330 .iter()
331 .map(|item| String::from(item.unwrap_or_default()))
332 .collect();
333
334 let stream_properties_jsonb = stream_properties_accessor.jsonb_value(ir)?;
336
337 let stream = Arc::new(StreamMetadata {
338 stream_id,
339 process_id: stream_process_id,
340 dependencies_metadata: ciborium::from_reader(dependencies_metadata)
341 .with_context(|| "decoding dependencies_metadata")?,
342 objects_metadata: ciborium::from_reader(objects_metadata)
343 .with_context(|| "decoding objects_metadata")?,
344 tags: stream_tags,
345 properties: Arc::new(stream_properties_jsonb),
346 });
347
348 let format = stream_format_column.value(ir)?.to_string();
349
350 if partition_nb_objects + block_nb_objects > config.max_nb_objects
352 && !partition_blocks.is_empty()
353 {
354 partitions.push(SourceDataBlocksInMemory {
356 blocks: partition_blocks,
357 block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
358 });
359 partition_blocks = vec![Arc::new(PartitionSourceBlock {
361 block,
362 stream: stream.clone(),
363 process: process.clone(),
364 format,
365 })];
366 partition_nb_objects = block_nb_objects;
367 } else {
368 partition_nb_objects += block_nb_objects;
370 partition_blocks.push(Arc::new(PartitionSourceBlock {
371 block,
372 stream: stream.clone(),
373 process: process.clone(),
374 format,
375 }));
376 }
377 }
378 }
379 if partition_nb_objects != 0 {
380 partitions.push(SourceDataBlocksInMemory {
381 blocks: partition_blocks,
382 block_ids_hash: partition_nb_objects.to_le_bytes().to_vec(),
383 });
384 }
385 Ok(partitions)
386}
387
388#[span_fn]
392pub async fn generate_process_jit_partitions(
393 config: &JitPartitionConfig,
394 lakehouse: Arc<LakehouseContext>,
395 blocks_view: &BlocksView,
396 query_time_range: &TimeRange,
397 process: Arc<ProcessMetadata>,
398 stream_tag: &str,
399) -> Result<Vec<SourceDataBlocksInMemory>> {
400 let part_provider = LivePartitionProvider::new(lakehouse.lake().db_pool.clone());
402 let view_set_name = blocks_view.get_view_set_name();
403 let view_instance_id = blocks_view.get_view_instance_id();
404 let partitions = instrument_named!(
405 part_provider.fetch(
406 &view_set_name,
407 &view_instance_id,
408 Some(*query_time_range),
409 blocks_view.get_file_schema_hash(),
410 ),
411 "live_partition_provider_fetch"
412 )
413 .await?;
414
415 let process_id = &process.process_id;
416 let begin_range_iso = query_time_range.begin.to_rfc3339();
417 let end_range_iso = query_time_range.end.to_rfc3339();
418 let sql = format!(
419 r#"SELECT MIN(insert_time) as min_insert_time, MAX(insert_time) as max_insert_time
420 FROM source
421 WHERE process_id = '{process_id}'
422 AND array_has( "streams.tags", '{stream_tag}' )
423 AND begin_time <= '{end_range_iso}'
424 AND end_time >= '{begin_range_iso}';"#
425 );
426
427 let reader_factory = lakehouse.reader_factory().clone();
428 let df = instrument_named!(
429 query_partitions(
430 lakehouse.runtime().clone(),
431 reader_factory,
432 lakehouse.lake().blob_storage.inner(),
433 blocks_view.get_file_schema(),
434 Arc::new(partitions),
435 &sql,
436 ),
437 "query_partitions"
438 )
439 .await?;
440 let rbs = instrument_named!(df.collect(), "collect_insert_time_range").await?;
441
442 if rbs.is_empty() || rbs[0].num_rows() == 0 {
443 return Ok(vec![]);
444 }
445
446 let min_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 0)?;
447 let max_insert_time = get_single_row_primitive_value::<TimestampNanosecondType>(&rbs, 1)?;
448
449 if min_insert_time == 0 || max_insert_time == 0 {
450 return Ok(vec![]);
451 }
452
453 let insert_time_range = TimeRange::new(
454 DateTime::from_timestamp_nanos(min_insert_time)
455 .duration_trunc(config.max_insert_time_slice)?,
456 DateTime::from_timestamp_nanos(max_insert_time)
457 .duration_trunc(config.max_insert_time_slice)?
458 + config.max_insert_time_slice,
459 );
460
461 let segment_source_partitions = instrument_named!(
462 PartitionCache::fetch_overlapping_insert_range_for_view(
463 &lakehouse.lake().db_pool,
464 blocks_view.get_view_set_name(),
465 blocks_view.get_view_instance_id(),
466 insert_time_range,
467 ),
468 "fetch_overlapping_insert_range_for_view"
469 )
470 .await?;
471
472 let mut begin_segment = insert_time_range.begin;
473 let mut end_segment = begin_segment + config.max_insert_time_slice;
474 let mut partitions = vec![];
475
476 while end_segment <= insert_time_range.end {
477 let insert_time_range = TimeRange::new(begin_segment, end_segment);
478 let mut segment_partitions = generate_process_jit_partitions_segment(
479 config,
480 lakehouse.clone(),
481 blocks_view,
482 &segment_source_partitions,
483 &insert_time_range,
484 process.clone(),
485 stream_tag,
486 )
487 .await?;
488 partitions.append(&mut segment_partitions);
489 begin_segment = end_segment;
490 end_segment = begin_segment + config.max_insert_time_slice;
491 }
492 Ok(partitions)
493}
494
495#[span_fn]
498pub async fn is_jit_partition_up_to_date(
499 pool: &sqlx::PgPool,
500 view_meta: ViewMetadata,
501 spec: &SourceDataBlocksInMemory,
502) -> Result<bool> {
503 let (min_insert_time, max_insert_time) =
504 get_part_insert_time_range(spec).with_context(|| "get_event_time_range")?;
505 let desc = format!(
506 "[{}, {}] {} {}",
507 min_insert_time.to_rfc3339(),
508 max_insert_time.to_rfc3339(),
509 *view_meta.view_set_name,
510 *view_meta.view_instance_id,
511 );
512
513 let rows = instrument_named!(
521 if min_insert_time == max_insert_time {
522 sqlx::query(
524 "SELECT file_schema_hash, source_data_hash
525 FROM lakehouse_partitions
526 WHERE view_set_name = $1
527 AND view_instance_id = $2
528 AND begin_insert_time = $3
529 AND end_insert_time = $3
530 ;",
531 )
532 .bind(&*view_meta.view_set_name)
533 .bind(&*view_meta.view_instance_id)
534 .bind(min_insert_time)
535 } else {
536 sqlx::query(
538 "SELECT file_schema_hash, source_data_hash
539 FROM lakehouse_partitions
540 WHERE view_set_name = $1
541 AND view_instance_id = $2
542 AND begin_insert_time <= $3
543 AND end_insert_time >= $4
544 ;",
545 )
546 .bind(&*view_meta.view_set_name)
547 .bind(&*view_meta.view_instance_id)
548 .bind(max_insert_time)
549 .bind(min_insert_time)
550 }
551 .fetch_all(pool),
552 "sql_select_matching_partitions"
553 )
554 .await
555 .with_context(|| "fetching matching partitions")?;
556 if rows.len() != 1 {
557 debug!("{desc}: found {} partitions (expected 1)", rows.len());
558 for (i, row) in rows.iter().enumerate() {
559 let part_file_schema: Vec<u8> = row.try_get("file_schema_hash")?;
560 let part_source_data: Vec<u8> = row.try_get("source_data_hash")?;
561 let source_row_count = hash_to_object_count(&part_source_data)?;
562 debug!(
563 "{desc}: partition {}: file_schema_hash={:?}, source_rows={}",
564 i, part_file_schema, source_row_count
565 );
566 }
567 info!("{desc}: found {} partitions", rows.len());
568 return Ok(false);
569 }
570 let r = &rows[0];
571 let part_file_schema: Vec<u8> = r.try_get("file_schema_hash")?;
572 if part_file_schema != view_meta.file_schema_hash {
573 warn!("{desc}: found matching partition with different file schema");
576 return Ok(false);
577 }
578 let part_source_data: Vec<u8> = r.try_get("source_data_hash")?;
579 let existing_count = hash_to_object_count(&part_source_data)?;
580 let required_count = hash_to_object_count(&spec.block_ids_hash)?;
581 if existing_count < required_count {
582 info!("{desc}: existing partition lacks source data: creating a new partition");
583 return Ok(false);
584 }
585 info!("{desc}: partition up to date");
586 Ok(true)
587}
588
589fn get_part_insert_time_range(
592 spec: &SourceDataBlocksInMemory,
593) -> Result<(DateTime<Utc>, DateTime<Utc>)> {
594 if spec.blocks.is_empty() {
595 anyhow::bail!("empty partition should not exist");
596 }
597 let min_insert_time = spec.blocks[0].block.insert_time;
599 let max_insert_time = spec.blocks[spec.blocks.len() - 1].block.insert_time;
600 Ok((min_insert_time, max_insert_time))
601}
602
603#[span_fn]
608pub async fn write_partition_from_blocks(
609 lake: Arc<DataLakeConnection>,
610 view_metadata: ViewMetadata,
611 schema: Arc<Schema>,
612 source_data: SourceDataBlocksInMemory,
613 block_processors: Arc<BlockProcessorMap>,
614) -> Result<()> {
615 if source_data.blocks.is_empty() {
616 anyhow::bail!("empty partition spec");
617 }
618 let min_insert_time = source_data.blocks[0].block.insert_time;
620 let max_insert_time = source_data.blocks[source_data.blocks.len() - 1]
621 .block
622 .insert_time;
623 let block_spec = BlockPartitionSpec {
624 view_metadata,
625 schema,
626 insert_range: TimeRange::new(min_insert_time, max_insert_time),
627 source_data: Arc::new(source_data),
628 block_processors,
629 };
630 let null_response_writer = Arc::new(ResponseWriter::new(None));
631 block_spec
632 .write(lake, null_response_writer)
633 .await
634 .with_context(|| "block_spec.write")?;
635 Ok(())
636}