From af5c138241f2bbb66733f88bd4b0ff17149f75a8 Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Mon, 28 Sep 2026 21:41:48 +0800 Subject: [PATCH 1/2] fix(datafusion): map column-less MERGE INSERT VALUES by position MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `WHEN NOT MATCHED THEN INSERT VALUES (s.a, s.b, ...)` without an explicit column list built an empty column->expression map, so `insert_select_clause` emitted `NULL` for every target column and dropped the source values. The merge then inserted all-NULL rows (or failed on a non-null column) instead of the source data — silent corruption for a column-less positional INSERT, which is valid standard SQL (the MERGE INSERT column list is optional) and accepted by DataFusion's parser. Map the VALUES to the table's columns by position when no column list is given, and reject a value/column count mismatch. Explicit-column and `INSERT *` clauses are unchanged. Added an end-to-end MERGE test using a column-less INSERT; it inserts all-NULL rows before the fix. --- .../integrations/datafusion/src/merge_into.rs | 116 ++++++++++++++---- 1 file changed, 94 insertions(+), 22 deletions(-) diff --git a/crates/integrations/datafusion/src/merge_into.rs b/crates/integrations/datafusion/src/merge_into.rs index 0bb9a3b22..74cc6e4e1 100644 --- a/crates/integrations/datafusion/src/merge_into.rs +++ b/crates/integrations/datafusion/src/merge_into.rs @@ -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?; @@ -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 { 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::>() - .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::>() + .join(", ")) } /// Parsed WHEN NOT MATCHED THEN INSERT clause. @@ -2191,6 +2215,54 @@ 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_update_and_insert() { let (_tmp, sql_context, table) = setup_append_only_table("t_cow_upsert").await; From 8178f94100eb4acbe175962ab56fd8e5b87ca21c Mon Sep 17 00:00:00 2001 From: jackylee-ch Date: Fri, 2 Oct 2026 09:38:47 +0800 Subject: [PATCH 2/2] fix(datafusion): validate positional MERGE INSERT arity upfront The value-count check for a column-less positional INSERT lived in `insert_select_clause`, which only runs while building unmatched-row batches. A fully matched CoW merge (and the data-evolution empty-batch path) skips that step, so an invalid `INSERT VALUES (s.a, s.b)` against a three-column target was accepted and the matched UPDATE committed. Move the arity check into `validate_merge_insert_columns`, which both execute paths call before any batch is built. `INSERT *`/`INSERT ROW` (no columns and no values) is unaffected. Added a regression: an all-matched MERGE with a short positional INSERT now fails without changing the target. --- .../integrations/datafusion/src/merge_into.rs | 71 +++++++++++++++++++ 1 file changed, 71 insertions(+) diff --git a/crates/integrations/datafusion/src/merge_into.rs b/crates/integrations/datafusion/src/merge_into.rs index 74cc6e4e1..15e8c82e8 100644 --- a/crates/integrations/datafusion/src/merge_into.rs +++ b/crates/integrations/datafusion/src/merge_into.rs @@ -1102,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(()) @@ -2263,6 +2281,59 @@ mod tests { ); } + #[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;