From e90637d0266ac08e15ce6eccb09bbdb20469458c Mon Sep 17 00:00:00 2001 From: Swanand Mulay <73115739+swanandx@users.noreply.github.com> Date: Thu, 23 Jul 2026 02:13:35 +0530 Subject: [PATCH] connector: postgres sink respect column DEFAULTs for omitted columns The INSERT used `SELECT *` from jsonb_populate_recordset, which named every column and fed NULL into columns absent from the Feldera schema, breaking DEFAULT values of columns. List only the supplied columns so Postgres applies each omitted column's DEFAULT. Signed-off-by: Swanand Mulay <73115739+swanandx@users.noreply.github.com> --- .../postgres/prepared_statements.rs | 102 ++++++++++++++- .../adapters/src/integrated/postgres/test.rs | 119 ++++++++++++++++++ 2 files changed, 219 insertions(+), 2 deletions(-) diff --git a/crates/adapters/src/integrated/postgres/prepared_statements.rs b/crates/adapters/src/integrated/postgres/prepared_statements.rs index a6435f122cd..6dc15ce5feb 100644 --- a/crates/adapters/src/integrated/postgres/prepared_statements.rs +++ b/crates/adapters/src/integrated/postgres/prepared_statements.rs @@ -22,13 +22,29 @@ impl RawQueries { .map(|f| f.name.sql_name()) .collect(); + // List supplied columns explicitly so Postgres applies DEFAULTs to any + // omitted column, instead of SELECT * forcing NULL into them (#6694). + let mut insert_columns: Vec = value_schema + .fields + .iter() + .map(|f| f.name.sql_name()) + .chain(config.extra_columns.iter().map(|k| format!(r#""{k}""#))) + .collect(); + + // CDC rows also carry the op/ts metadata columns. + if matches!(config.mode, PostgresWriteMode::Cdc) { + insert_columns.push(format!(r#""{}""#, config.cdc_op_column)); + insert_columns.push(format!(r#""{}""#, config.cdc_ts_column)); + } + let insert_columns = insert_columns.join(", "); + let mut raw_queries = RawQueries::default(); match config.mode { PostgresWriteMode::Cdc => { // In CDC mode, everything is an INSERT into the event log raw_queries.insert = format!( - r#"INSERT INTO "{table}" SELECT * FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb)"#, + r#"INSERT INTO "{table}" ({insert_columns}) SELECT {insert_columns} FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb)"#, ); // For CDC mode, upsert and delete operations also use INSERT raw_queries.upsert = raw_queries.insert.clone(); @@ -59,7 +75,7 @@ impl RawQueries { }; raw_queries.insert = format!( - r#"INSERT INTO "{table}" SELECT * FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb) ON CONFLICT {on_conflict}"#, + r#"INSERT INTO "{table}" ({insert_columns}) SELECT {insert_columns} FROM jsonb_populate_recordset(NULL::"{table}", $1::jsonb) ON CONFLICT {on_conflict}"#, ); } } @@ -169,3 +185,85 @@ impl PreparedStatements { }) } } + +#[cfg(test)] +mod tests { + use super::*; + use feldera_types::program_schema::{ColumnType, Field}; + use std::collections::BTreeMap; + + fn relation(name: &str, columns: &[&str]) -> Relation { + Relation::new( + name.into(), + columns + .iter() + .map(|c| Field::new((*c).into(), ColumnType::varchar(true))) + .collect(), + false, + BTreeMap::new(), + ) + } + + fn writer_config(mode: PostgresWriteMode, extra_columns: &[&str]) -> PostgresWriterConfig { + serde_json::from_value(serde_json::json!({ + "uri": "postgres://localhost", + "table": "t", + "mode": mode.to_string(), + "extra_columns": extra_columns, + })) + .unwrap() + } + + // INSERT lists supplied columns, never SELECT * (#6694). + #[test] + fn materialized_insert_lists_supplied_columns() { + let key = relation("k", &["id"]); + let value = relation("v", &["id", "name"]); + let queries = RawQueries::new( + &key, + &value, + &writer_config(PostgresWriteMode::Materialized, &["audit"]), + ); + + assert!(!queries.insert.contains("SELECT *"), "{}", queries.insert); + assert!(queries.insert.contains(r#"(id, name, "audit")"#)); + assert!(queries.insert.contains(r#"SELECT id, name, "audit" FROM"#)); + } + + // Case-sensitive columns stay quoted. + #[test] + fn materialized_insert_quotes_case_sensitive_columns() { + let key = relation("k", &[r#""Id""#]); + let value = relation("v", &[r#""Id""#, r#""Name""#]); + let queries = RawQueries::new( + &key, + &value, + &writer_config(PostgresWriteMode::Materialized, &[]), + ); + + assert!(queries.insert.contains(r#"("Id", "Name")"#)); + assert!(queries.insert.contains(r#"SELECT "Id", "Name" FROM"#)); + } + + // CDC INSERT also lists the op/ts metadata columns. + #[test] + fn cdc_insert_lists_supplied_and_metadata_columns() { + let key = relation("k", &["id"]); + let value = relation("v", &["id", "name"]); + let queries = RawQueries::new(&key, &value, &writer_config(PostgresWriteMode::Cdc, &[])); + + assert!(!queries.insert.contains("SELECT *")); + assert!( + queries + .insert + .contains(r#"(id, name, "__feldera_op", "__feldera_ts")"#) + ); + assert!( + queries + .insert + .contains(r#"SELECT id, name, "__feldera_op", "__feldera_ts" FROM"#) + ); + assert_eq!(queries.insert, queries.upsert); + assert_eq!(queries.insert, queries.delete); + } +} diff --git a/crates/adapters/src/integrated/postgres/test.rs b/crates/adapters/src/integrated/postgres/test.rs index 533cf65759c..fbd1eabab5d 100644 --- a/crates/adapters/src/integrated/postgres/test.rs +++ b/crates/adapters/src/integrated/postgres/test.rs @@ -275,9 +275,27 @@ mod pg { impl TempPgTable { fn new(name: &str, uri: String, pk: bool, tls: &Option) -> Self { + Self::new_with_extra_columns(name, uri, pk, tls, "") + } + + /// Same as `new`, but appends `extra_columns` to the table: column + /// definitions absent from the Feldera schema, e.g. + /// `"served_at TIMESTAMP DEFAULT now() NOT NULL"`. + fn new_with_extra_columns( + name: &str, + uri: String, + pk: bool, + tls: &Option, + extra_columns: &str, + ) -> Self { let mut client = pg_connect(&uri, tls); let pk = if pk { "PRIMARY KEY" } else { "" }; + let extra_columns = if extra_columns.is_empty() { + String::new() + } else { + format!(",\n {extra_columns}") + }; client .execute( @@ -322,6 +340,7 @@ CREATE TABLE {name} ( map_ JSONB, __feldera_op CHAR, __feldera_ts BIGINT + {extra_columns} )"# ), &[], @@ -490,6 +509,18 @@ CREATE TABLE {name} ( TempPgTable::new(name, uri, pk, tls) } + /// Create the table with extra column definitions absent from the + /// Feldera schema. + pub fn create_table_with_extra_columns( + name: &str, + uri: String, + pk: bool, + tls: &Option, + extra_columns: &str, + ) -> TempPgTable { + TempPgTable::new_with_extra_columns(name, uri, pk, tls, extra_columns) + } + pub fn test_circuit( config: PipelineConfig, ) -> (Controller, crossbeam::channel::Receiver) { @@ -950,6 +981,94 @@ fn pg_insert0(mode: PostgresWriteMode) { .expect("timeout: failed to insert data into postgres"); } +// A NOT NULL DEFAULT column (served_at) omitted from the Feldera schema must be +// populated from its default, not NULL (#6694). +#[test] +#[serial] +fn test_pg_insert_omitted_default_column() { + let table_name = unique_pg_name("test_pg_default"); + let url = postgres_url(); + let verify_url = url.clone(); + let timeout_ms = 60_000; + + let mut data: Vec = (0..100).map(|_| rand::random()).collect(); + + let mut temp_file = NamedTempFile::new().unwrap(); + for datum in data.iter() { + let mut serializer = serde_json::Serializer::new(Vec::new()); + datum + .serialize_with_context(&mut serializer, &SqlSerdeConfig::default()) + .unwrap(); + temp_file + .as_file_mut() + .write_all(&serializer.into_inner()) + .unwrap(); + temp_file.write_all(b"\n").unwrap(); + } + + let config = serde_json::from_value(json!({ + "name": "test", + "workers": 4, + "inputs": { + "ins": { + "stream": "test_input1", + "transport": { "name": "file_input", "config": { "path": temp_file.path() } }, + "format": { "name": "json", "config": { "update_format": "raw", "array": false } } + } + }, + "outputs": { + "test_output1": { + "stream": "test_output1", + "transport": { "name": "postgres_output", "config": { "uri": url, "table": &table_name } }, + "index": "idx" + } + } + })) + .unwrap(); + + // The Feldera view (PostgresTestStruct) has no served_at column. + let mut table = PostgresTestStruct::create_table_with_extra_columns( + &table_name, + url, + true, + &None, + "served_at TIMESTAMP DEFAULT now() NOT NULL", + ); + + let (controller, err_receiver) = PostgresTestStruct::test_circuit(config); + controller.start(); + data.sort(); + + wait( + || { + let mut got = table + .query() + .into_iter() + .map(PostgresTestStruct::from) + .collect::>(); + got.sort(); + data == got || !err_receiver.is_empty() + }, + timeout_ms, + ) + .expect("timeout: failed to insert data into postgres"); + + assert!(err_receiver.is_empty(), "connector reported an error"); + + // Every row's served_at was filled by the DEFAULT rather than left NULL. + let mut client = pg::pg_connect(&verify_url, &None); + let rows = client + .query(&format!("SELECT served_at FROM {table_name}"), &[]) + .expect("failed to query table"); + assert_eq!(rows.len(), data.len()); + assert!( + rows.iter().all(|r| r + .get::<_, Option>("served_at") + .is_some()), + "served_at should be populated by its DEFAULT" + ); +} + #[test] #[serial] fn test_pg_insert() {