micromegas_analytics/lakehouse/
merge.rs1use super::{
2 lakehouse_context::LakehouseContext,
3 partition::Partition,
4 partition_cache::PartitionCache,
5 partition_source_data::hash_to_object_count,
6 partitioned_execution_plan::OrderingBounds,
7 partitioned_table_provider::PartitionedTableProvider,
8 query::make_session_context,
9 session_configurator::SessionConfigurator,
10 view::{ScanSortColumn, View},
11 view_factory::ViewFactory,
12 write_partition::{PartitionRowSet, write_partition_from_rows},
13};
14use crate::{response_writer::Logger, time::TimeRange};
15use anyhow::{Context, Result};
16use async_trait::async_trait;
17use datafusion::{
18 arrow::datatypes::Schema,
19 execution::SendableRecordBatchStream,
20 physical_plan::{displayable, execute_stream},
21 prelude::*,
22 sql::TableReference,
23};
24use futures::stream::StreamExt;
25use micromegas_tracing::prelude::*;
26use std::fmt::Debug;
27use std::sync::Arc;
28use xxhash_rust::xxh32::xxh32;
29
30pub struct MergeQueryResult {
32 pub stream: SendableRecordBatchStream,
34 pub ordering_honored: bool,
40}
41
42#[async_trait]
44pub trait PartitionMerger: Send + Sync + Debug {
45 async fn execute_merge_query(
47 &self,
48 lakehouse: Arc<LakehouseContext>,
49 partitions_to_merge: Arc<Vec<Partition>>,
50 partitions_all_views: Arc<PartitionCache>,
51 insert_range: TimeRange,
52 ) -> Result<MergeQueryResult>;
53}
54
55#[derive(Debug)]
57pub struct QueryMerger {
58 view_factory: Arc<ViewFactory>,
59 session_configurator: Arc<dyn SessionConfigurator>,
60 file_schema: Arc<Schema>,
61 query: Arc<String>,
62 merge_scan_ordering: Vec<ScanSortColumn>,
63}
64
65impl QueryMerger {
66 pub fn new(
67 view_factory: Arc<ViewFactory>,
68 session_configurator: Arc<dyn SessionConfigurator>,
69 file_schema: Arc<Schema>,
70 query: Arc<String>,
71 ) -> Self {
72 Self {
73 view_factory,
74 session_configurator,
75 file_schema,
76 query,
77 merge_scan_ordering: vec![],
78 }
79 }
80
81 pub fn with_merge_scan_ordering(mut self, ordering: Vec<ScanSortColumn>) -> Self {
86 self.merge_scan_ordering = ordering;
87 self
88 }
89}
90
91#[async_trait]
92impl PartitionMerger for QueryMerger {
93 async fn execute_merge_query(
94 &self,
95 lakehouse: Arc<LakehouseContext>,
96 partitions_to_merge: Arc<Vec<Partition>>,
97 partitions_all_views: Arc<PartitionCache>,
98 insert_range: TimeRange,
99 ) -> Result<MergeQueryResult> {
100 let reader_factory = lakehouse.reader_factory().clone();
101 let ctx = make_session_context(
102 lakehouse.clone(),
103 partitions_all_views,
104 Some(insert_range),
105 self.view_factory.clone(),
106 self.session_configurator.clone(),
107 true,
108 )
109 .await?;
110 let src_table = PartitionedTableProvider::with_ordering(
111 self.file_schema.clone(),
112 reader_factory,
113 partitions_to_merge,
114 self.merge_scan_ordering.clone(),
115 OrderingBounds::InsertTime,
116 );
117 ctx.register_table(
118 TableReference::Bare {
119 table: "source".into(),
120 },
121 Arc::new(src_table),
122 )?;
123
124 if self.merge_scan_ordering.is_empty() {
125 let stream = ctx
126 .sql(&self.query)
127 .await?
128 .execute_stream()
129 .await
130 .with_context(|| "merged_df.execute_stream")?;
131 return Ok(MergeQueryResult {
132 stream,
133 ordering_honored: true,
134 });
135 }
136
137 ctx.state_ref()
142 .write()
143 .config_mut()
144 .options_mut()
145 .optimizer
146 .repartition_file_scans = false;
147
148 let df = ctx.sql(&self.query).await?;
149 let task_ctx = Arc::new(df.task_ctx());
150 let plan = df
151 .create_physical_plan()
152 .await
153 .with_context(|| "creating physical plan for merge query")?;
154
155 let partition_count = plan.properties().output_partitioning().partition_count();
156 if partition_count != 1 {
157 anyhow::bail!(
158 "merge query {:?} (insert_range=[{}, {}]) produced a {partition_count}-partition \
159 physical plan; executing it would coalesce partitions and destroy the declared \
160 ordering. This likely means repartition_file_scans did not take effect.",
161 self.query,
162 insert_range.begin.to_rfc3339(),
163 insert_range.end.to_rfc3339()
164 );
165 }
166
167 let plan_str = displayable(plan.as_ref()).indent(true).to_string();
168 let ordering_honored =
169 !plan_str.contains("SortExec") && !plan_str.contains("SortPreservingMergeExec");
170 if !ordering_honored {
171 warn!(
172 "merge query {:?} (insert_range=[{}, {}]) did not elide its declared ordering -- \
173 the merge will still produce a correctly ordered result, but it will buffer in \
174 memory instead of streaming. Plan:\n{plan_str}",
175 self.query,
176 insert_range.begin.to_rfc3339(),
177 insert_range.end.to_rfc3339()
178 );
179 }
180
181 let stream =
182 execute_stream(plan, task_ctx).with_context(|| "executing merge query plan")?;
183 Ok(MergeQueryResult {
184 stream,
185 ordering_honored,
186 })
187 }
188}
189
190fn partition_set_stats(
191 view: Arc<dyn View>,
192 filtered_partitions: &[Partition],
193) -> Result<(i64, i64)> {
194 let mut sum_size: i64 = 0;
195 let mut source_hash: i64 = 0;
196 let latest_file_schema_hash = view.get_file_schema_hash();
197 for p in filtered_partitions {
198 source_hash = if p.source_data_hash.len() == std::mem::size_of::<i64>() {
201 source_hash + hash_to_object_count(&p.source_data_hash)?
202 } else {
203 xxh32(&p.source_data_hash, source_hash as u32).into()
205 };
206
207 sum_size += p.file_size;
208
209 if p.view_metadata.file_schema_hash != latest_file_schema_hash {
210 anyhow::bail!(
211 "incompatible file schema with [{},{}]",
212 p.begin_insert_time().to_rfc3339(),
213 p.end_insert_time().to_rfc3339()
214 );
215 }
216 }
217 Ok((sum_size, source_hash))
218}
219
220pub async fn create_merged_partition(
222 partitions_to_merge: Arc<PartitionCache>,
223 partitions_all_views: Arc<PartitionCache>,
224 lakehouse: Arc<LakehouseContext>,
225 view: Arc<dyn View>,
226 insert_range: TimeRange,
227 logger: Arc<dyn Logger>,
228) -> Result<()> {
229 let view_set_name = &view.get_view_set_name();
230 let view_instance_id = &view.get_view_instance_id();
231 let desc = format!(
232 "[{}, {}] {view_set_name} {view_instance_id}",
233 insert_range.begin.to_rfc3339(),
234 insert_range.end.to_rfc3339()
235 );
236 let mut filtered_partitions = partitions_to_merge
239 .filter_inside_range(view_set_name, view_instance_id, insert_range)
240 .partitions;
241 if filtered_partitions.len() != partitions_to_merge.len() {
242 warn!("partitions_to_merge was not filtered properly");
243 }
244 if filtered_partitions.len() < 2 {
245 logger
246 .write_log_entry(format!("{desc}: not enough partitions to merge"))
247 .await
248 .with_context(|| "writing log")?;
249 return Ok(());
250 }
251 let (sum_size, source_hash) = partition_set_stats(view.clone(), &filtered_partitions)
252 .with_context(|| "partition_set_stats")?;
253 logger
254 .write_log_entry(format!(
255 "{desc}: merging {} partitions sum_size={sum_size}",
256 filtered_partitions.len()
257 ))
258 .await
259 .with_context(|| "write_log_entry")?;
260 filtered_partitions.sort_by_key(|p| p.begin_insert_time());
261 let merged_sort_order = view.get_merged_partition_sort_order(&filtered_partitions);
263 let merge_result = view
264 .merge_partitions(
265 lakehouse.clone(),
266 Arc::new(filtered_partitions),
267 partitions_all_views,
268 insert_range,
269 )
270 .await
271 .with_context(|| "view.merge_partitions")?;
272 let mut merged_stream = merge_result.stream;
275 let (tx, rx) = tokio::sync::mpsc::channel(1);
276 let view_copy = view.clone();
277 let lake = lakehouse.lake().clone();
278 let join_handle = spawn_with_context(write_partition_from_rows(
279 lake,
280 view_copy.get_meta(),
281 view_copy.get_file_schema(),
282 insert_range,
283 source_hash.to_le_bytes().to_vec(),
284 merged_sort_order,
285 rx,
286 logger.clone(),
287 ));
288 let compute_time_bounds = view.get_time_bounds();
289 let ctx =
290 SessionContext::new_with_config_rt(SessionConfig::default(), lakehouse.runtime().clone());
291 let stream_result: Result<()> = async {
292 while let Some(rb_res) = merged_stream.next().await {
293 let rb = rb_res.with_context(|| "receiving record_batch from stream")?;
294 let event_time_range = compute_time_bounds
295 .get_time_bounds(ctx.read_batch(rb.clone()).with_context(|| "read_batch")?)
296 .await?;
297 tx.send(Ok(PartitionRowSet::new(event_time_range, rb)))
298 .await
299 .with_context(|| "sending partition row set")?;
300 }
301 Ok(())
302 }
303 .await;
304
305 match stream_result {
306 Ok(()) => {
307 drop(tx);
308 join_handle.await??;
309 Ok(())
310 }
311 Err(e) => {
312 warn!("aborting merge partition write for {desc}: {e:?}");
313 let _ = tx.send(Err(anyhow::anyhow!("merge stream aborted"))).await;
314 drop(tx);
315 match join_handle.await {
316 Ok(Ok(())) => {}
317 Ok(Err(writer_err)) => {
318 debug!("merge writer task error during abort: {writer_err:?}");
319 }
320 Err(join_err) => {
321 warn!("merge writer task panicked during abort: {join_err:?}");
322 }
323 }
324 Err(e)
325 }
326 }
327}