Skip to main content

micromegas_analytics/lakehouse/
batch_update.rs

1use super::{
2    lakehouse_context::LakehouseContext, merge::create_merged_partition,
3    partition_cache::PartitionCache, partition_source_data::hash_to_object_count, view::View,
4};
5use crate::{response_writer::Logger, time::TimeRange};
6use anyhow::{Context, Result};
7use chrono::TimeDelta;
8use micromegas_tracing::prelude::*;
9use std::sync::Arc;
10
11/// Defines the strategy for creating a new partition.
12pub enum PartitionCreationStrategy {
13    /// Create the partition from the source data.
14    CreateFromSource,
15    /// Merge existing partitions.
16    MergeExisting(Arc<PartitionCache>),
17    /// Abort the partition creation.
18    Abort,
19}
20
21// verify_overlapping_partitions returns true to continue and make a new partition,
22// returns false to abort (existing partition is up to date or there is a problem)
23async fn verify_overlapping_partitions(
24    existing_partitions_all_views: &PartitionCache,
25    insert_range: TimeRange,
26    view_set_name: &str,
27    view_instance_id: &str,
28    file_schema_hash: &[u8],
29    source_data_hash: &[u8],
30    logger: Arc<dyn Logger>,
31) -> Result<PartitionCreationStrategy> {
32    let desc = format!(
33        "[{}, {}] {view_set_name} {view_instance_id}",
34        insert_range.begin.to_rfc3339(),
35        insert_range.end.to_rfc3339()
36    );
37    if source_data_hash.len() != std::mem::size_of::<i64>() {
38        anyhow::bail!("Source data hash should be a i64");
39    }
40    let nb_source_events = hash_to_object_count(source_data_hash)?;
41    let filtered = existing_partitions_all_views.filter(
42        view_set_name,
43        view_instance_id,
44        file_schema_hash,
45        insert_range,
46    );
47    if filtered.partitions.is_empty() {
48        logger
49            .write_log_entry(format!("{desc}: matching partitions not found"))
50            .await?;
51        return Ok(PartitionCreationStrategy::CreateFromSource);
52    }
53    let mut existing_source_hash: i64 = 0;
54    let nb_existing_partitions = filtered.partitions.len();
55    for part in &filtered.partitions {
56        let begin = part.begin_insert_time();
57        let end = part.end_insert_time();
58        if begin < insert_range.begin || end > insert_range.end {
59            logger
60                .write_log_entry(format!(
61                    "{desc}: found overlapping partition [{}, {}], aborting the update",
62                    begin.to_rfc3339(),
63                    end.to_rfc3339()
64                ))
65                .await?;
66            return Ok(PartitionCreationStrategy::Abort);
67        }
68        if part.source_data_hash.len() == std::mem::size_of::<i64>() {
69            existing_source_hash += hash_to_object_count(&part.source_data_hash)?
70        } else {
71            // old hash that does not represent the number of events
72            logger
73                .write_log_entry(format!(
74                    "{desc}: found partition with incompatible source hash: recreate"
75                ))
76                .await?;
77            return Ok(PartitionCreationStrategy::CreateFromSource);
78        }
79    }
80
81    if nb_source_events != existing_source_hash {
82        logger
83            .write_log_entry(format!(
84                "{desc}: existing partitions do not match source data ({nb_source_events} vs {existing_source_hash}) : creating a new partition"
85            ))
86            .await?;
87        return Ok(PartitionCreationStrategy::CreateFromSource);
88    }
89
90    if nb_existing_partitions > 1 {
91        return Ok(PartitionCreationStrategy::MergeExisting(Arc::new(filtered)));
92    }
93
94    logger
95        .write_log_entry(format!(
96            "{desc}: already up to date, nb_source_events={nb_source_events}"
97        ))
98        .await?;
99    Ok(PartitionCreationStrategy::Abort)
100}
101
102/// Re-checks the same partial-overlap condition `verify_overlapping_partitions` guards, without
103/// its source-hash freshness comparison (the "already up to date" `Abort`, which regeneration
104/// exists to bypass). Used only by `regenerate_partition`.
105///
106/// Filters the existing-partitions snapshot hash-agnostically (`filter_insert_range` +
107/// view/instance match) rather than via `PartitionCache::filter`'s hash-exact match: a
108/// partially-overlapping partition written under an older schema hash must still be caught here,
109/// since `retire_partitions`'s range-containment delete is equally hash-agnostic and would not
110/// retire it either.
111fn verify_force_regeneration_alignment(
112    existing_partitions_all_views: &PartitionCache,
113    insert_range: TimeRange,
114    view_set_name: &str,
115    view_instance_id: &str,
116) -> Result<()> {
117    let filtered = existing_partitions_all_views.filter_insert_range(insert_range);
118    for part in &filtered.partitions {
119        if *part.view_metadata.view_set_name != view_set_name
120            || *part.view_metadata.view_instance_id != view_instance_id
121        {
122            continue;
123        }
124        let begin = part.begin_insert_time();
125        let end = part.end_insert_time();
126        if begin < insert_range.begin || end > insert_range.end {
127            anyhow::bail!(
128                "regenerate_partitions: regeneration bucket [{}, {}] does not fully contain \
129                 existing partition [{}, {}] for {view_set_name}/{view_instance_id} -- \
130                 the requested range/delta must exactly cover the partition(s) being regenerated",
131                insert_range.begin.to_rfc3339(),
132                insert_range.end.to_rfc3339(),
133                begin.to_rfc3339(),
134                end.to_rfc3339(),
135            );
136        }
137    }
138    Ok(())
139}
140
141#[span_fn]
142async fn materialize_partition(
143    existing_partitions_all_views: Arc<PartitionCache>,
144    lakehouse: Arc<LakehouseContext>,
145    insert_range: TimeRange,
146    view: Arc<dyn View>,
147    logger: Arc<dyn Logger>,
148) -> Result<()> {
149    let view_set_name = view.get_view_set_name();
150    let partition_spec = view
151        .make_batch_partition_spec(
152            lakehouse.clone(),
153            existing_partitions_all_views.clone(),
154            insert_range,
155        )
156        .await
157        .with_context(|| "make_batch_partition_spec")?;
158    // Allow empty partition specs to be written - write_partition_from_rows
159    // will create an empty partition record
160    let view_instance_id = view.get_view_instance_id();
161    let strategy = verify_overlapping_partitions(
162        &existing_partitions_all_views,
163        insert_range,
164        &view_set_name,
165        &view_instance_id,
166        &view.get_file_schema_hash(),
167        &partition_spec.get_source_data_hash(),
168        logger.clone(),
169    )
170    .await
171    .with_context(|| "verify_overlapping_partitions")?;
172    if let PartitionCreationStrategy::Abort = &strategy {
173        return Ok(());
174    }
175
176    let new_delta = view.get_max_partition_time_delta(&strategy);
177    if new_delta < (insert_range.end - insert_range.begin) {
178        if let PartitionCreationStrategy::MergeExisting(partition_cache) = &strategy
179            && partition_cache
180                .partitions
181                .iter()
182                .all(|p| (p.end_insert_time() - p.begin_insert_time()) == new_delta)
183        {
184            let desc = format!(
185                "[{}, {}] {view_set_name} {view_instance_id}",
186                insert_range.begin.to_rfc3339(),
187                insert_range.end.to_rfc3339()
188            );
189            logger
190                .write_log_entry(format!("{desc}: subpartitions already present",))
191                .await?;
192            return Ok(());
193        }
194
195        return Box::pin(materialize_partition_range(
196            existing_partitions_all_views,
197            lakehouse.clone(),
198            view,
199            insert_range,
200            new_delta,
201            logger,
202        ))
203        .await
204        .with_context(|| "materialize_partition_range");
205    }
206
207    match strategy {
208        PartitionCreationStrategy::CreateFromSource => {
209            partition_spec
210                .write(lakehouse.lake().clone(), logger)
211                .await
212                .with_context(|| "writing partition")?;
213        }
214        PartitionCreationStrategy::MergeExisting(partitions_to_merge) => {
215            create_merged_partition(
216                partitions_to_merge,
217                existing_partitions_all_views,
218                lakehouse,
219                view,
220                insert_range,
221                logger,
222            )
223            .await
224            .with_context(|| "create_merged_partition")?;
225        }
226        PartitionCreationStrategy::Abort => {}
227    }
228
229    Ok(())
230}
231
232/// Materializes partitions within a given time range.
233#[span_fn]
234pub async fn materialize_partition_range(
235    existing_partitions_all_views: Arc<PartitionCache>,
236    lakehouse: Arc<LakehouseContext>,
237    view: Arc<dyn View>,
238    insert_range: TimeRange,
239    partition_time_delta: TimeDelta,
240    logger: Arc<dyn Logger>,
241) -> Result<()> {
242    let mut begin_part = insert_range.begin;
243    let mut end_part = begin_part + partition_time_delta;
244    while end_part <= insert_range.end {
245        let partition_insert_range = TimeRange::new(begin_part, end_part);
246        let insert_time_filtered =
247            Arc::new(existing_partitions_all_views.filter_insert_range(partition_insert_range));
248        materialize_partition(
249            insert_time_filtered,
250            lakehouse.clone(),
251            partition_insert_range,
252            view.clone(),
253            logger.clone(),
254        )
255        .await
256        .with_context(|| "materialize_partition")?;
257        begin_part = end_part;
258        end_part = begin_part + partition_time_delta;
259    }
260    Ok(())
261}
262
263/// Regenerates the partition(s) covering `insert_range` directly from source data, bypassing the
264/// "already up to date" freshness check `materialize_partition_range` would otherwise stop at. See
265/// `tasks/blocks_view_ordered_merges_plan.md`'s Design ยง3.
266///
267/// **Invariant callers must uphold**: `(begin, end, delta)` must exactly cover the boundaries of
268/// the partition(s) being regenerated -- a misaligned range/delta means the new partition's range
269/// does not fully contain the old one, so `retire_partitions` never retires it, leaving silent
270/// duplicate rows. This is enforced by validating that `delta` exactly tiles `(begin, end)` before
271/// any partition is written, and (per bucket) by `verify_force_regeneration_alignment`.
272///
273/// Both checks read a snapshot and are advisory UX: the authoritative guard is the
274/// `lakehouse_partitions_no_overlap` exclusion constraint, which makes the insert fail loudly if
275/// a conflicting partition was committed concurrently (e.g. by the maintenance daemon merging
276/// buckets after the snapshot was taken).
277#[span_fn]
278pub async fn regenerate_partition_range(
279    existing_partitions_all_views: Arc<PartitionCache>,
280    lakehouse: Arc<LakehouseContext>,
281    view: Arc<dyn View>,
282    insert_range: TimeRange,
283    partition_time_delta: TimeDelta,
284    logger: Arc<dyn Logger>,
285) -> Result<()> {
286    // chrono::TimeDelta implements no Rem/% operator, so tile-checking is done on nanoseconds.
287    let span = (insert_range.end - insert_range.begin)
288        .num_nanoseconds()
289        .expect("time range span should fit in an i64 number of nanoseconds");
290    let step = partition_time_delta
291        .num_nanoseconds()
292        .expect("partition_time_delta should fit in an i64 number of nanoseconds");
293    if !(step > 0 && span >= step && span % step == 0) {
294        anyhow::bail!(
295            "regenerate_partitions: delta ({partition_time_delta}) does not exactly tile the \
296             requested range [{}, {}] -- range/delta must exactly cover the partition(s) being \
297             regenerated",
298            insert_range.begin.to_rfc3339(),
299            insert_range.end.to_rfc3339(),
300        );
301    }
302    let mut begin_part = insert_range.begin;
303    let mut end_part = begin_part + partition_time_delta;
304    while end_part <= insert_range.end {
305        let bucket = TimeRange::new(begin_part, end_part);
306        let filtered = Arc::new(existing_partitions_all_views.filter_insert_range(bucket));
307        regenerate_partition(
308            filtered,
309            lakehouse.clone(),
310            bucket,
311            view.clone(),
312            logger.clone(),
313        )
314        .await
315        .with_context(|| "regenerate_partition")?;
316        begin_part = end_part;
317        end_part = begin_part + partition_time_delta;
318    }
319    Ok(())
320}
321
322/// Regenerates a single bucket from source, replacing whatever aligned partition(s) currently
323/// cover it. Unlike `materialize_partition` it always writes from source -- never merges, never
324/// aborts on freshness -- and never subdivides: the bucket is exactly one partition. The
325/// transactional retire+insert in the write path replaces the old partition(s) atomically, so a
326/// failure rolls back and leaves the existing partition untouched.
327#[span_fn]
328async fn regenerate_partition(
329    existing_partitions_all_views: Arc<PartitionCache>,
330    lakehouse: Arc<LakehouseContext>,
331    insert_range: TimeRange,
332    view: Arc<dyn View>,
333    logger: Arc<dyn Logger>,
334) -> Result<()> {
335    let view_set_name = view.get_view_set_name();
336    let view_instance_id = view.get_view_instance_id();
337    verify_force_regeneration_alignment(
338        &existing_partitions_all_views,
339        insert_range,
340        &view_set_name,
341        &view_instance_id,
342    )
343    .with_context(|| "verify_force_regeneration_alignment")?;
344    let partition_spec = view
345        .make_batch_partition_spec(
346            lakehouse.clone(),
347            existing_partitions_all_views,
348            insert_range,
349        )
350        .await
351        .with_context(|| "make_batch_partition_spec")?;
352    partition_spec
353        .write(lakehouse.lake().clone(), logger)
354        .await
355        .with_context(|| "writing partition")
356}