Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
187 changes: 165 additions & 22 deletions crates/integrations/datafusion/src/merge_into.rs
Original file line number Diff line number Diff line change
Expand Up @@ -907,7 +907,7 @@ async fn build_insert_batches_inner(
format!(" WHERE {}", conditions.join(" AND "))
};

let select_clause = insert_select_clause(ins, table_fields);
let select_clause = insert_select_clause(ins, table_fields)?;
let sql = format!("SELECT {select_clause} FROM {tmp_name} AS {s_alias}{where_clause}");

let batches = ctx.ctx().sql(&sql).await?.collect().await?;
Expand Down Expand Up @@ -995,32 +995,56 @@ fn strip_non_source_columns(
///
/// When the INSERT specifies explicit columns (`INSERT (col2, col1) VALUES (expr2, expr1)`),
/// the output must be reordered to match the table schema so that `write_arrow_batch`
/// (which reads columns by positional index) maps them correctly.
fn insert_select_clause(ins: &MergeInsertClause, table_fields: &[DataField]) -> String {
/// (which reads columns by positional index) maps them correctly. When the INSERT omits
/// the column list (`INSERT VALUES (expr1, expr2, ...)`), the values map to the table's
/// columns by position, following standard SQL, where the MERGE INSERT column list is
/// optional (unlike Spark, whose grammar requires an explicit list or `INSERT *`).
fn insert_select_clause(ins: &MergeInsertClause, table_fields: &[DataField]) -> DFResult<String> {
if ins.columns.is_empty() && ins.value_exprs.is_empty() {
"*".to_string()
} else {
// Build column_name -> expression mapping from the INSERT clause
let col_expr_map: HashMap<&str, &str> = ins
.columns
.iter()
.zip(ins.value_exprs.iter())
.map(|(col, expr)| (col.as_str(), expr.as_str()))
.collect();
return Ok("*".to_string());
}

// Emit SELECT in table schema order
table_fields
if ins.columns.is_empty() {
// No explicit column list: the values fill the table's columns by
// position. Without this branch the empty column->expr map below would
// emit `NULL` for every column and silently drop the source values,
// inserting all-NULL rows (or failing on a non-null column).
if ins.value_exprs.len() != table_fields.len() {
return Err(DataFusionError::Plan(format!(
"MERGE INSERT has {} value(s) but the table has {} column(s); \
list the target columns explicitly or provide one value per column",
ins.value_exprs.len(),
table_fields.len()
)));
}
return Ok(table_fields
.iter()
.map(|field| {
match col_expr_map.get(field.name()) {
Some(expr) => format!("{expr} AS {}", quote_identifier(field.name())),
// Column not in INSERT list — fill with NULL
None => format!("NULL AS {}", quote_identifier(field.name())),
}
})
.zip(ins.value_exprs.iter())
.map(|(field, expr)| format!("{expr} AS {}", quote_identifier(field.name())))
.collect::<Vec<_>>()
.join(", ")
.join(", "));
}

// Build column_name -> expression mapping from the INSERT clause
let col_expr_map: HashMap<&str, &str> = ins
.columns
.iter()
.zip(ins.value_exprs.iter())
.map(|(col, expr)| (col.as_str(), expr.as_str()))
.collect();

// Emit SELECT in table schema order
Ok(table_fields
.iter()
.map(|field| {
match col_expr_map.get(field.name()) {
Some(expr) => format!("{expr} AS {}", quote_identifier(field.name())),
// Column not in INSERT list — fill with NULL
None => format!("NULL AS {}", quote_identifier(field.name())),
}
})
.collect::<Vec<_>>()
.join(", "))
}

/// Parsed WHEN NOT MATCHED THEN INSERT clause.
Expand Down Expand Up @@ -1078,6 +1102,24 @@ fn validate_merge_insert_columns(
) -> DFResult<()> {
for insert in inserts {
validate_target_columns(&insert.columns, table_fields, "MERGE INSERT")?;
// A column-less positional INSERT maps its VALUES to the table columns by
// position, so the value count must match the column count. Validate it
// here, upfront: this runs in both CoW and data-evolution modes before any
// batch is built, whereas `insert_select_clause` only runs when unmatched
// rows exist — an all-matched MERGE would otherwise skip the arity check
// and commit. `INSERT *`/`INSERT ROW` (no columns and no values) is left
// untouched.
if insert.columns.is_empty()
&& !insert.value_exprs.is_empty()
&& insert.value_exprs.len() != table_fields.len()
{
return Err(DataFusionError::Plan(format!(
"MERGE INSERT has {} value(s) but the table has {} column(s); \
list the target columns explicitly or provide one value per column",
insert.value_exprs.len(),
table_fields.len()
)));
}
}

Ok(())
Expand Down Expand Up @@ -2191,6 +2233,107 @@ mod tests {
);
}

#[tokio::test]
async fn test_cow_merge_insert_not_matched_without_columns() {
let (_tmp, sql_context, table) = setup_append_only_table("t_cow_ins_nocols").await;

sql_context
.sql("CREATE TABLE paimon.test_db.source (id INT, name VARCHAR, value INT)")
.await
.unwrap();
sql_context
.sql("INSERT INTO paimon.test_db.source VALUES (4, 'dave', 40), (5, 'eve', 50)")
.await
.unwrap()
.collect()
.await
.unwrap();

// No column list: the VALUES map to the table's columns by position.
// Before the fix this dropped the source values and inserted all-NULL rows.
let merge = parse_merge(
"MERGE INTO paimon.test_db.t_cow_ins_nocols t USING paimon.test_db.source s \
ON t.id = s.id \
WHEN NOT MATCHED THEN INSERT VALUES (s.id, s.name, s.value)",
);
execute_merge_into(&sql_context, &merge, table, true)
.await
.unwrap();

let batches = sql_context
.sql("SELECT id, name, value FROM paimon.test_db.t_cow_ins_nocols ORDER BY id")
.await
.unwrap()
.collect()
.await
.unwrap();

let rows = collect_rows(&batches);
assert_eq!(
rows,
vec![
(1, "alice".to_string(), 10),
(2, "bob".to_string(), 20),
(3, "charlie".to_string(), 30),
(4, "dave".to_string(), 40),
(5, "eve".to_string(), 50),
]
);
}

#[tokio::test]
async fn test_cow_merge_insert_arity_validated_even_when_all_matched() {
let (_tmp, sql_context, table) = setup_append_only_table("t_cow_ins_arity").await;

sql_context
.sql("CREATE TABLE paimon.test_db.source (id INT, name VARCHAR, value INT)")
.await
.unwrap();
// id = 1 matches the target, so the NOT MATCHED INSERT branch is never
// built — the arity check must still run upfront.
sql_context
.sql("INSERT INTO paimon.test_db.source VALUES (1, 'ALICE', 99)")
.await
.unwrap()
.collect()
.await
.unwrap();

// The column-less INSERT lists 2 values for a 3-column table. Even with no
// row to insert, the invalid arity must be rejected before execution so the
// matched UPDATE is not applied.
let merge = parse_merge(
"MERGE INTO paimon.test_db.t_cow_ins_arity t USING paimon.test_db.source s \
ON t.id = s.id \
WHEN MATCHED THEN UPDATE SET value = s.value \
WHEN NOT MATCHED THEN INSERT VALUES (s.id, s.name)",
);
let err = execute_merge_into(&sql_context, &merge, table, true)
.await
.unwrap_err();
assert!(
err.to_string().contains("value(s) but the table has"),
"got {err}"
);

// The target is unchanged: the matched UPDATE did not apply.
let batches = sql_context
.sql("SELECT id, name, value FROM paimon.test_db.t_cow_ins_arity ORDER BY id")
.await
.unwrap()
.collect()
.await
.unwrap();
assert_eq!(
collect_rows(&batches),
vec![
(1, "alice".to_string(), 10),
(2, "bob".to_string(), 20),
(3, "charlie".to_string(), 30),
]
);
}

#[tokio::test]
async fn test_cow_merge_update_and_insert() {
let (_tmp, sql_context, table) = setup_append_only_table("t_cow_upsert").await;
Expand Down
Loading