laminar_connectors/lakehouse/delta_io/
merge.rs1use super::{debug, info, ConnectorError, DeltaTable, DeltaWriteAttemptError, RecordBatch};
4
5#[cfg(feature = "delta-lake")]
7#[derive(Debug)]
8pub struct MergeResult {
9 pub rows_inserted: usize,
11 pub rows_updated: usize,
13 pub rows_deleted: usize,
15}
16
17#[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 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 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 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}