diff --git a/crates/integrations/datafusion/src/merge_into.rs b/crates/integrations/datafusion/src/merge_into.rs index 0bb9a3b22..15e8c82e8 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. @@ -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(()) @@ -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;