Skip to main content

micromegas_analytics/lakehouse/
merge.rs

1use 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
30/// The outcome of running a merge query.
31pub struct MergeQueryResult {
32    /// The merged rows.
33    pub stream: SendableRecordBatchStream,
34    /// Whether the merger's declared scan ordering (if any) was honored by the physical plan
35    /// without falling back to a buffering `Sort`/`SortPreservingMerge` node. Always `true` when
36    /// no ordering was declared to DataFusion in the first place -- it is only ever computed
37    /// dynamically by an ordering-declaring `QueryMerger`. This drives only a memory-regression
38    /// warning; it never gates the recorded `sort_order` (see `View::get_merged_partition_sort_order`).
39    pub ordering_honored: bool,
40}
41
42/// A trait for merging partitions.
43#[async_trait]
44pub trait PartitionMerger: Send + Sync + Debug {
45    /// Executes the merge query.
46    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/// A `PartitionMerger` that executes a SQL query to merge partitions.
56#[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    /// Declares an ordering the merge's source scan already satisfies (see
82    /// `PartitionedTableProvider::with_ordering`), letting DataFusion elide the merge query's
83    /// `Sort` node instead of buffering. Default: empty (no declared ordering, matching today's
84    /// behavior for every existing caller).
85    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        // Ordering-declared merge: force the source scan into a single sequential file group
138        // (Design §1 point 3) so the declared ordering can be elided instead of re-sorted, then
139        // build the physical plan once, inspect it, and execute that exact plan -- never
140        // planning or building twice.
141        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        // for some time all the hashes will actually be the number of events in the source data
199        // when views have different hash algos, we should delegate to the view the creation of the merged hash
200        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            //previous hash algo
204            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
220/// Creates a merged partition from a set of existing partitions.
221pub 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    // we are not looking for intersecting partitions, but only those that fit completely in the range
237    // otherwise we'd get duplicated records
238    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    // Computed before merge_partitions runs: a pure function of the input slice alone (Design §4).
262    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    // A defeated ordering elision (ordering_honored: false) is already warned about, with the
273    // offending plan, inside QueryMerger::execute_merge_query.
274    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}