Skip to main content

laminar_connectors/lakehouse/delta_io/
merge.rs

1//! Changelog MERGE execution and result accounting.
2
3use super::{debug, info, ConnectorError, DeltaTable, DeltaWriteAttemptError, RecordBatch};
4
5/// Result of a MERGE (upsert) operation.
6#[cfg(feature = "delta-lake")]
7#[derive(Debug)]
8pub struct MergeResult {
9    /// Number of rows inserted.
10    pub rows_inserted: usize,
11    /// Number of rows updated.
12    pub rows_updated: usize,
13    /// Number of rows deleted.
14    pub rows_deleted: usize,
15}
16
17/// Atomic changelog MERGE: inserts, updates, and deletes in one Delta commit.
18///
19/// The source batch must contain an `_op` column (Utf8) with values:
20/// - `"I"`, `"U"`, `"r"` → upsert (update if matched, insert if not)
21/// - `"D"` → delete matched rows
22///
23/// Columns prefixed with `_` are excluded from SET clauses but remain
24/// in the source `DataFrame` for predicate filtering.
25///
26/// # Errors
27///
28/// Returns `ConnectorError::WriteError` if the merge fails.
29#[cfg(feature = "delta-lake")]
30#[allow(clippy::too_many_lines)]
31pub(crate) async fn merge_changelog(
32    table: DeltaTable,
33    source_batch: RecordBatch,
34    key_columns: &[String],
35    schema_evolution: bool,
36    writer_properties: Option<deltalake::parquet::file::properties::WriterProperties>,
37    ctx: &datafusion::prelude::SessionContext,
38) -> Result<(DeltaTable, MergeResult), DeltaWriteAttemptError> {
39    use datafusion::prelude::*;
40    use deltalake::kernel::transaction::CommitProperties;
41
42    const CDC_COLUMNS: &[&str] = &["_op", "_ts_ms"];
43
44    if source_batch.num_rows() == 0 {
45        return Ok((
46            table,
47            MergeResult {
48                rows_inserted: 0,
49                rows_updated: 0,
50                rows_deleted: 0,
51            },
52        ));
53    }
54
55    debug!(
56        key_columns = ?key_columns,
57        source_rows = source_batch.num_rows(),
58        "performing atomic changelog MERGE"
59    );
60
61    let source_df = ctx.read_batch(source_batch).map_err(|e| {
62        ConnectorError::WriteError(format!("failed to create source DataFrame: {e}"))
63    })?;
64
65    // Join predicate: target.k1 = source.k1 AND ...
66    let predicate = key_columns
67        .iter()
68        .map(|k| col(format!("target.{k}")).eq(col(format!("source.{k}"))))
69        .reduce(Expr::and)
70        .ok_or_else(|| {
71            ConnectorError::ConfigurationError("merge requires at least one key column".into())
72        })?;
73
74    let source_schema = source_df.schema().clone();
75    let key_set: std::collections::HashSet<&str> = key_columns.iter().map(String::as_str).collect();
76
77    // Exclude CDC metadata columns from SET clauses (preserve user columns like _id).
78    let all_user_columns: Vec<String> = source_schema
79        .fields()
80        .iter()
81        .map(|f| f.name().clone())
82        .filter(|name| !CDC_COLUMNS.contains(&name.as_str()))
83        .collect();
84
85    let non_key_user_columns: Vec<String> = all_user_columns
86        .iter()
87        .filter(|c| !key_set.contains(c.as_str()))
88        .cloned()
89        .collect();
90
91    // Predicates for conditional clause execution.
92    let upsert_pred = col("source._op").in_list(vec![lit("I"), lit("U"), lit("r")], false);
93    let delete_pred = col("source._op").eq(lit("D"));
94
95    let non_key_for_update = non_key_user_columns;
96    let all_for_insert = all_user_columns;
97
98    let mut merge_builder = table
99        .merge(source_df, predicate)
100        .with_source_alias("source")
101        .with_target_alias("target")
102        .with_commit_properties(CommitProperties::default())
103        .when_matched_update(|update| {
104            let mut u = update.predicate(upsert_pred.clone());
105            for col_name in &non_key_for_update {
106                u = u.update(col_name.as_str(), col(format!("source.{col_name}")));
107            }
108            u
109        })
110        .map_err(|e| ConnectorError::WriteError(format!("merge matched-update failed: {e}")))?
111        .when_matched_delete(|delete| delete.predicate(delete_pred))
112        .map_err(|e| ConnectorError::WriteError(format!("merge matched-delete failed: {e}")))?
113        .when_not_matched_insert(|insert| {
114            let mut ins = insert.predicate(upsert_pred);
115            for col_name in &all_for_insert {
116                ins = ins.set(col_name.as_str(), col(format!("source.{col_name}")));
117            }
118            ins
119        })
120        .map_err(|e| ConnectorError::WriteError(format!("merge not-matched-insert failed: {e}")))?;
121
122    if schema_evolution {
123        merge_builder = merge_builder.with_merge_schema(true);
124    }
125
126    if let Some(props) = writer_properties {
127        merge_builder = merge_builder.with_writer_properties(props);
128    }
129
130    let (table, metrics) = merge_builder.await.map_err(DeltaWriteAttemptError::Delta)?;
131
132    let result = MergeResult {
133        rows_inserted: metrics.num_target_rows_inserted,
134        rows_updated: metrics.num_target_rows_updated,
135        rows_deleted: metrics.num_target_rows_deleted,
136    };
137
138    info!(
139        rows_inserted = result.rows_inserted,
140        rows_updated = result.rows_updated,
141        rows_deleted = result.rows_deleted,
142        "Delta Lake changelog MERGE complete"
143    );
144
145    Ok((table, result))
146}