micromegas_analytics/lakehouse/
batch_update.rs1use 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
11pub enum PartitionCreationStrategy {
13 CreateFromSource,
15 MergeExisting(Arc<PartitionCache>),
17 Abort,
19}
20
21async 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 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
102fn 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 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#[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#[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 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#[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}