From 0eae8c91f2c07f6aef53decae147dc66b42fa416 Mon Sep 17 00:00:00 2001 From: Hasyimi Bahrudin Date: Mon, 5 Oct 2026 20:11:10 +0800 Subject: [PATCH] fix: columns tracked by name instead of Iceberg field id MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - A restart registered each table with its discovered schema, whose field ids are column positions: after a column was dropped, the ones behind it took others' ids and were written under the wrong Iceberg column. An existing table now keeps its Iceberg schema; the stream's Relation messages evolve it, in order. - A column dropped and re-added under the same name kept the old field id, so the dropped values came back. A re-add — seen as the column returning after it was seen dropped, or out of the order Postgres appends re-added columns in — now renames the old column to `__dropped_`, values intact, and adds a new column. - pg2iceberg's own Parquet reader (compaction, TOAST resolution, FileIndex rebuild, verify) matched columns by name. It now matches by field id, as Iceberg does, and reads files written before an `int`→`long` / `float`→`double` promotion as the new type. - `oid` maps to `long` (it's unsigned 32-bit; values past 2^31 wrapped negative). Existing `int` columns are promoted by the next Relation. Rows staged before a re-add but materialized after it still land in the new column (schema changes apply ahead of rows already staged); pinned as a failing DST test for a follow-up that applies schema changes in log order. Co-Authored-By: Claude Opus 5.5 --- crates/pg2iceberg-core/src/typemap.rs | 18 +- crates/pg2iceberg-iceberg/src/lib.rs | 278 ++++++++++++++++++ crates/pg2iceberg-iceberg/src/reader.rs | 143 +++++++-- crates/pg2iceberg-logical/src/materializer.rs | 206 +++++-------- crates/pg2iceberg-pg/src/prod/value_decode.rs | 8 +- crates/pg2iceberg-sim/src/postgres.rs | 16 + crates/pg2iceberg-tests/tests/dst.rs | 54 ++++ 7 files changed, 559 insertions(+), 164 deletions(-) diff --git a/crates/pg2iceberg-core/src/typemap.rs b/crates/pg2iceberg-core/src/typemap.rs index c519f2f..786de59 100644 --- a/crates/pg2iceberg-core/src/typemap.rs +++ b/crates/pg2iceberg-core/src/typemap.rs @@ -87,8 +87,9 @@ pub const DEFAULT_NUMERIC_SCALE: u8 = 18; pub fn map_pg_to_iceberg(pg: PgType) -> Result { let iceberg = match pg { PgType::Bool => IcebergType::Boolean, - PgType::Int2 | PgType::Int4 | PgType::Oid => IcebergType::Int, - PgType::Int8 => IcebergType::Long, + PgType::Int2 | PgType::Int4 => IcebergType::Int, + // `oid` is unsigned 32-bit: values past 2^31 don't fit `int`. + PgType::Int8 | PgType::Oid => IcebergType::Long, PgType::Float4 => IcebergType::Float, PgType::Float8 => IcebergType::Double, PgType::Numeric { precision, scale } => return map_numeric(precision, scale), @@ -194,18 +195,17 @@ mod tests { } #[test] - fn int2_int4_oid_all_map_to_int() { - for ty in [PgType::Int2, PgType::Int4, PgType::Oid] { + fn int2_int4_map_to_int() { + for ty in [PgType::Int2, PgType::Int4] { assert_eq!(map_pg_to_iceberg(ty).unwrap().iceberg, IcebergType::Int); } } #[test] - fn int8_maps_to_long() { - assert_eq!( - map_pg_to_iceberg(PgType::Int8).unwrap().iceberg, - IcebergType::Long - ); + fn int8_and_oid_map_to_long() { + for ty in [PgType::Int8, PgType::Oid] { + assert_eq!(map_pg_to_iceberg(ty).unwrap().iceberg, IcebergType::Long); + } } #[test] diff --git a/crates/pg2iceberg-iceberg/src/lib.rs b/crates/pg2iceberg-iceberg/src/lib.rs index 5b0751c..969e905 100644 --- a/crates/pg2iceberg-iceberg/src/lib.rs +++ b/crates/pg2iceberg-iceberg/src/lib.rs @@ -152,6 +152,18 @@ pub enum SchemaChange { name: String, new_ty: pg2iceberg_core::IcebergType, }, + /// Rename a column, keeping its field id — and so its data. Moves a + /// dropped column out of the way of a re-added one with the same + /// name (see [`dropped_column_name`]). + RenameColumn { from: String, to: String }, +} + +/// The name a dropped column takes when the source re-adds a column with +/// its name: the re-added column is new (Postgres gives it no values), so +/// it gets a new field id, and the dropped one's values stay readable +/// under this name. +pub fn dropped_column_name(name: &str, field_id: i32) -> String { + format!("{name}__dropped_{field_id}") } /// True if `to` is a spec-legal Iceberg promotion of `from`. Reference: @@ -252,11 +264,145 @@ pub fn apply_schema_changes( } col.ty = *new_ty; } + SchemaChange::RenameColumn { from, to } => { + if schema.columns.iter().any(|c| c.name == *to) { + return Err(IcebergError::Conflict(format!( + "RenameColumn: column {to} already exists" + ))); + } + let col = schema + .columns + .iter_mut() + .find(|c| c.name == *from) + .ok_or_else(|| { + IcebergError::NotFound(format!("RenameColumn: column {from} not in schema")) + })?; + col.name = to.clone(); + } } } Ok(()) } +/// The changes that bring `table` (the Iceberg schema) in line with the +/// source table's current columns, `source` (name and type, in the +/// source's column order). +/// +/// Columns match by name — Iceberg field ids stay with their columns — +/// except one the source dropped and re-added: Postgres gives that a new, +/// empty column, so it becomes a new Iceberg column (new field id) and +/// the dropped one is renamed out of its way ([`dropped_column_name`]), +/// its values intact. A re-add shows either way: +/// +/// - the column was seen dropped: it's in `dropped`; +/// - it's out of order: Postgres appends a re-added column after every +/// surviving one and never reorders columns otherwise, while `table` +/// keeps the order columns were added in. So from the first column +/// that comes before one it used to follow, the rest were re-added. +/// +/// A column re-added while it was already last, with no change to the +/// table in between, can't be told from one that was never dropped. +/// +/// Columns the source no longer has are soft-dropped (kept, nullable); +/// type changes must be legal promotions, and never of a key column. +pub fn reconcile_columns( + table: &TableSchema, + source: &[(String, pg2iceberg_core::IcebergType)], + dropped: &std::collections::BTreeSet, +) -> Result> { + use std::collections::{BTreeMap, BTreeSet}; + let position: BTreeMap<&str, (usize, &pg2iceberg_core::ColumnSchema)> = table + .columns + .iter() + .enumerate() + .map(|(i, c)| (c.name.as_str(), (i, c))) + .collect(); + + let mut readded: BTreeSet<&str> = BTreeSet::new(); + let mut last_kept: Option = None; + let mut out_of_order = false; + for (name, _) in source { + let Some(&(at, col)) = position.get(name.as_str()) else { + continue; + }; + if col.is_primary_key { + continue; + } + out_of_order |= last_kept.is_some_and(|last| at < last); + if out_of_order || dropped.contains(name) { + readded.insert(name); + } else { + last_kept = Some(at); + } + } + + let mut changes = Vec::new(); + for (name, ty) in source { + match position.get(name.as_str()) { + Some(&(_, col)) if readded.contains(name.as_str()) => { + if !col.nullable { + changes.push(SchemaChange::DropColumn { name: name.clone() }); + } + changes.push(SchemaChange::RenameColumn { + from: name.clone(), + to: dropped_column_name(name, col.field_id), + }); + changes.push(SchemaChange::AddColumn { + name: name.clone(), + ty: *ty, + nullable: true, + }); + } + Some(&(_, col)) => { + if col.ty == *ty { + continue; + } + // Key columns are part of the equality-delete + // predicate; promoting one would invalidate every prior + // delete file's keys. + if col.is_primary_key { + return Err(IcebergError::Other(format!( + "cannot promote primary-key column {name}: {:?} → {ty:?} \ + (PK type is part of the equality-delete contract; \ + changing it requires a full re-snapshot)", + col.ty + ))); + } + if !is_legal_type_promotion(col.ty, *ty) { + return Err(IcebergError::Other(format!( + "column {name} type change {:?} → {ty:?} is not a legal Iceberg \ + promotion. Allowed: int→long, float→double, decimal \ + precision increase. Other changes (narrowing, cross-family) \ + require a full re-snapshot.", + col.ty + ))); + } + changes.push(SchemaChange::PromoteColumnType { + name: name.clone(), + new_ty: *ty, + }); + } + // Iceberg requires an added column to be optional, so files + // written before it read it as NULL. + None => changes.push(SchemaChange::AddColumn { + name: name.clone(), + ty: *ty, + nullable: true, + }), + } + } + for col in &table.columns { + // Key columns are never dropped; a dropped column already + // nullable needs nothing. + if !col.is_primary_key && !col.nullable && !source.iter().any(|(n, _)| *n == col.name) { + changes.push(SchemaChange::DropColumn { + name: col.name.clone(), + }); + } + } + Ok(changes) +} + #[async_trait] pub trait Catalog: Send + Sync { async fn ensure_namespace(&self, ns: &Namespace) -> Result<()>; @@ -546,3 +692,135 @@ mod schema_change_tests { assert_eq!(amount.field_id, 3); } } + +#[cfg(test)] +mod reconcile_tests { + use super::*; + use pg2iceberg_core::{ColumnSchema, IcebergType as T}; + use std::collections::BTreeSet; + + /// `id` (key), `note`, `qty` — field ids 1, 2, 3. + fn table() -> TableSchema { + let col = |name: &str, field_id: i32, ty: T, key: bool| ColumnSchema { + name: name.into(), + field_id, + ty, + nullable: !key, + is_primary_key: key, + }; + TableSchema { + ident: TableIdent { + namespace: Namespace(vec!["public".into()]), + name: "t".into(), + }, + columns: vec![ + col("id", 1, T::Int, true), + ColumnSchema { + nullable: false, + ..col("note", 2, T::String, false) + }, + col("qty", 3, T::Int, false), + ], + partition_spec: vec![], + pg_schema: None, + } + } + + fn source(cols: &[(&str, T)]) -> Vec<(String, T)> { + cols.iter().map(|(n, t)| (n.to_string(), *t)).collect() + } + + /// `(name, field id, nullable)` of the schema after `changes`. + fn after(changes: &[SchemaChange]) -> Vec<(String, i32, bool)> { + let mut s = table(); + apply_schema_changes(&mut s, changes).unwrap(); + s.columns + .into_iter() + .map(|c| (c.name, c.field_id, c.nullable)) + .collect() + } + + fn cols(v: &[(&str, i32, bool)]) -> Vec<(String, i32, bool)> { + v.iter().map(|(n, i, b)| (n.to_string(), *i, *b)).collect() + } + + #[test] + fn unchanged_columns_need_nothing() { + let src = source(&[("id", T::Int), ("note", T::String), ("qty", T::Int)]); + assert!(reconcile_columns(&table(), &src, &BTreeSet::new()) + .unwrap() + .is_empty()); + } + + #[test] + fn columns_keep_their_field_ids_after_one_is_dropped() { + // Discovery after `DROP COLUMN note` numbers qty 2 by position; + // it keeps field id 3, and note stays, soft-dropped. + let src = source(&[("id", T::Int), ("qty", T::Int)]); + let changes = reconcile_columns(&table(), &src, &BTreeSet::new()).unwrap(); + assert_eq!( + after(&changes), + cols(&[("id", 1, false), ("note", 2, true), ("qty", 3, true)]) + ); + } + + #[test] + fn a_new_column_gets_a_new_field_id() { + let src = source(&[ + ("id", T::Int), + ("note", T::String), + ("qty", T::Int), + ("tag", T::String), + ]); + let changes = reconcile_columns(&table(), &src, &BTreeSet::new()).unwrap(); + assert_eq!(after(&changes).last(), Some(&("tag".into(), 4, true))); + } + + #[test] + fn a_column_re_added_out_of_order_is_new() { + // `DROP COLUMN note; ADD COLUMN note`: Postgres appends it. + let src = source(&[("id", T::Int), ("qty", T::Int), ("note", T::String)]); + let changes = reconcile_columns(&table(), &src, &BTreeSet::new()).unwrap(); + assert_eq!( + after(&changes), + cols(&[ + ("id", 1, false), + ("note__dropped_2", 2, true), + ("qty", 3, true), + ("note", 4, true), + ]) + ); + } + + #[test] + fn a_column_seen_dropped_and_back_is_new() { + // qty was last, so its re-add keeps the order; it was seen gone. + let src = source(&[("id", T::Int), ("note", T::String), ("qty", T::Int)]); + let dropped = BTreeSet::from(["qty".to_string()]); + let changes = reconcile_columns(&table(), &src, &dropped).unwrap(); + assert_eq!( + after(&changes), + cols(&[ + ("id", 1, false), + ("note", 2, false), + ("qty__dropped_3", 3, true), + ("qty", 4, true), + ]) + ); + } + + #[test] + fn type_changes_must_be_legal_promotions_of_non_key_columns() { + let promote = source(&[("id", T::Int), ("note", T::String), ("qty", T::Long)]); + let changes = reconcile_columns(&table(), &promote, &BTreeSet::new()).unwrap(); + assert!(matches!( + changes.as_slice(), + [SchemaChange::PromoteColumnType { name, new_ty: T::Long }] if name == "qty" + )); + let narrow = source(&[("id", T::Int), ("note", T::Int), ("qty", T::Int)]); + assert!(reconcile_columns(&table(), &narrow, &BTreeSet::new()).is_err()); + let key = source(&[("id", T::Long), ("note", T::String), ("qty", T::Int)]); + let err = reconcile_columns(&table(), &key, &BTreeSet::new()).unwrap_err(); + assert!(err.to_string().contains("primary-key")); + } +} diff --git a/crates/pg2iceberg-iceberg/src/reader.rs b/crates/pg2iceberg-iceberg/src/reader.rs index 0bae7d4..bb307f2 100644 --- a/crates/pg2iceberg-iceberg/src/reader.rs +++ b/crates/pg2iceberg-iceberg/src/reader.rs @@ -2,8 +2,16 @@ //! //! Inverse of [`crate::writer::TableWriter::prepare`]'s data-file output. //! Used by TOAST resolution (re-fetching prior column values for an -//! `UPDATE ... SET col = unchanged_toast`) and by the `verify` -//! subcommand. +//! `UPDATE ... SET col = unchanged_toast`), compaction, FileIndex +//! rebuilds and the `verify` subcommand. +//! +//! Like any Iceberg reader, it matches a schema column to the file's +//! column with the same field id, not the same name: a column renamed +//! since the file was written keeps its values, and a column re-added +//! under an old name doesn't inherit them. A file written before a +//! column was added reads it as NULL, and one written before a type +//! promotion (`int` → `long`, `float` → `double`) reads as the promoted +//! type. use crate::writer::WriterError; use arrow_array::RecordBatch; @@ -35,8 +43,8 @@ pub fn read_data_file(bytes: &[u8], cols: &[ColumnSchema]) -> Result> { /// Decode a Parquet data file the way the Iceberg spec reads one: /// each of `cols` is matched to the file's column carrying the same /// `PARQUET:field_id`, not the same name, and a column the file doesn't -/// have reads as `NULL`. A query engine reading the table resolves -/// columns this way; [`read_data_file`] matches by name. +/// have reads as `NULL`. Deliberately separate from [`RowBatches`] — +/// tests use it as an independent oracle for that reader. pub fn read_data_file_by_field_id(bytes: Bytes, cols: &[ColumnSchema]) -> Result> { let reader = ParquetRecordBatchReaderBuilder::try_new(bytes) .map_err(|e| WriterError::Encode(format!("parquet reader: {e}")))? @@ -93,8 +101,14 @@ impl RowBatches { .map_err(|e| WriterError::Encode(format!("parquet reader: {e}")))?; let parquet_schema = builder.parquet_schema(); let wanted = (0..parquet_schema.num_columns()).filter(|&i| { - let name = parquet_schema.column(i).path().string(); - cols.iter().any(|c| c.name == name) + let column = parquet_schema.column(i); + let info = column.self_type().get_basic_info(); + if info.has_id() { + cols.iter().any(|c| c.field_id == info.id()) + } else { + let name = column.path().string(); + cols.iter().any(|c| c.name == name) + } }); let projection = ProjectionMask::leaves(parquet_schema, wanted); let reader = builder @@ -121,28 +135,57 @@ impl Iterator for RowBatches { } } +/// The batch's column for each of `cols`: by field id (by name for a +/// column the file wrote without one), `None` where the file has none. +fn resolve_columns<'a>( + batch: &'a RecordBatch, + cols: &[ColumnSchema], +) -> Vec> { + let schema = batch.schema(); + let mut by_id: BTreeMap = BTreeMap::new(); + let mut by_name: BTreeMap<&str, usize> = BTreeMap::new(); + for (i, f) in schema.fields().iter().enumerate() { + match f + .metadata() + .get("PARQUET:field_id") + .and_then(|id| id.parse().ok()) + { + Some(id) => { + by_id.insert(id, i); + } + None => { + by_name.insert(f.name().as_str(), i); + } + } + } + cols.iter() + .map(|c| { + by_id + .get(&c.field_id) + .or_else(|| by_name.get(c.name.as_str())) + .map(|&i| batch.column(i).as_ref()) + }) + .collect() +} + fn decode_batch(batch: &RecordBatch, cols: &[ColumnSchema]) -> Result> { + let arrays = resolve_columns(batch, cols); let mut out = Vec::with_capacity(batch.num_rows()); for i in 0..batch.num_rows() { let mut row: Row = BTreeMap::new(); - for col in cols { + for (col, arr) in cols.iter().zip(&arrays) { let key = ColumnName(col.name.clone()); - // Schema-evolution-tolerant read: if the parquet - // file pre-dates an `AddColumn`, the file simply - // doesn't carry the new column and Iceberg readers - // project NULL. Mirrors that behavior here so a - // mid-stream `ALTER TABLE ADD COLUMN` doesn't - // require backfilling old data files. Equality-delete - // reads (cols = PK-only) still hit every column - // they need, since PKs aren't dropped. - let Some(arr) = batch.column_by_name(&col.name) else { + // A file written before the column was added doesn't have + // it: Iceberg readers project NULL, so a mid-stream `ALTER + // TABLE ADD COLUMN` needs no backfill of old files. + let Some(arr) = *arr else { row.insert(key, PgValue::Null); continue; }; if arr.is_null(i) { row.insert(key, PgValue::Null); } else { - row.insert(key, decode_value(col.ty, arr.as_ref(), i)?); + row.insert(key, decode_value(col.ty, arr, i)?); } } out.push(row); @@ -165,9 +208,17 @@ fn decode_value(ty: IcebergType, arr: &dyn Array, i: usize) -> Result { Ok(match ty { IcebergType::Boolean => PgValue::Bool(cast!(BooleanArray)?.value(i)), IcebergType::Int => PgValue::Int4(cast!(Int32Array)?.value(i)), - IcebergType::Long => PgValue::Int8(cast!(Int64Array)?.value(i)), + // A file written before an `int` → `long` promotion. + IcebergType::Long => match arr.as_any().downcast_ref::() { + Some(old) => PgValue::Int8(old.value(i).into()), + None => PgValue::Int8(cast!(Int64Array)?.value(i)), + }, IcebergType::Float => PgValue::Float4(cast!(Float32Array)?.value(i)), - IcebergType::Double => PgValue::Float8(cast!(Float64Array)?.value(i)), + // A file written before a `float` → `double` promotion. + IcebergType::Double => match arr.as_any().downcast_ref::() { + Some(old) => PgValue::Float8(old.value(i).into()), + None => PgValue::Float8(cast!(Float64Array)?.value(i)), + }, IcebergType::String => PgValue::Text(cast!(StringArray)?.value(i).to_string()), IcebergType::Date => PgValue::Date(DaysSinceEpoch(cast!(Date32Array)?.value(i))), IcebergType::Timestamp => { @@ -297,6 +348,60 @@ mod tests { assert_eq!(decoded[1], row(2, 20)); } + /// A data file of `schema_id_qty()` holding `(1, 10)`. + fn id_qty_file() -> Bytes { + let rows = vec![MaterializedRow { + op: Op::Insert, + row: row(1, 10), + unchanged_cols: vec![], + unchanged_from: None, + }]; + let prepared = TableWriter::new(schema_id_qty()) + .prepare(&rows, &crate::FileIndex::new()) + .unwrap(); + prepared.data.into_iter().next().unwrap().chunk.bytes + } + + fn read(bytes: Bytes, cols: Vec) -> Row { + RowBatches::new(bytes, &cols, 16) + .unwrap() + .next() + .unwrap() + .unwrap() + .remove(0) + } + + #[test] + fn columns_resolve_by_field_id_not_name() { + let [id, qty] = schema_id_qty().columns.try_into().unwrap(); + // qty (field 2) since renamed, and a new column under its old + // name (field 3): the values follow the field id. + let renamed = ColumnSchema { + name: "qty__dropped_2".into(), + nullable: true, + ..qty.clone() + }; + let readded = ColumnSchema { + field_id: 3, + nullable: true, + ..qty + }; + let got = read(id_qty_file(), vec![id, renamed, readded]); + assert_eq!(got[&col("qty__dropped_2")], PgValue::Int4(10)); + assert_eq!(got[&col("qty")], PgValue::Null); + } + + #[test] + fn promoted_columns_read_as_the_new_type() { + let [id, qty] = schema_id_qty().columns.try_into().unwrap(); + let long = ColumnSchema { + ty: IcebergType::Long, + ..qty + }; + let got = read(id_qty_file(), vec![id, long]); + assert_eq!(got[&col("qty")], PgValue::Int8(10)); + } + #[test] fn round_trip_preserves_nulls() { let schema = TableSchema { diff --git a/crates/pg2iceberg-logical/src/materializer.rs b/crates/pg2iceberg-logical/src/materializer.rs index dbbbe5a..6ce1cac 100644 --- a/crates/pg2iceberg-logical/src/materializer.rs +++ b/crates/pg2iceberg-logical/src/materializer.rs @@ -36,9 +36,9 @@ use pg2iceberg_iceberg::meta::{ self as meta_schema, CheckpointStats, CompactionStats, FlushStats, MaintenanceStats, }; use pg2iceberg_iceberg::{ - fold_events, promote_re_inserts, read_data_file, rebuild_from_catalog, resolve_unchanged_cols, - Catalog, DataFile, FileIndex, IcebergError, MaterializedRow, PkKey, PreparedCommit, - TableWriter, WriterError, + fold_events, promote_re_inserts, read_data_file, rebuild_from_catalog, reconcile_columns, + resolve_unchanged_cols, Catalog, DataFile, FileIndex, IcebergError, MaterializedRow, PkKey, + PreparedCommit, TableWriter, WriterError, }; use pg2iceberg_stream::codec::decode_chunk; use pg2iceberg_stream::{BlobStore, MatEvent, StreamError}; @@ -378,6 +378,10 @@ struct TableEntry { /// [`Materializer::register_table`] entry, preserving existing /// behavior where the lifecycle gates the whole snapshot phase. gated_until_snapshot: bool, + /// Columns of `schema` the source was last seen without: soft-dropped. + /// One that comes back was re-added, so it's a new column (see + /// [`reconcile_columns`]). + dropped: BTreeSet, } /// Log entries fetched per `read_log` call while a cycle drains a table. @@ -832,9 +836,24 @@ impl Materializer { pub async fn register_table(&mut self, schema: TableSchema) -> Result<()> { let ident = schema.ident.clone(); self.catalog.ensure_namespace(&ident.namespace).await?; - if self.catalog.load_table(&ident).await?.is_none() { - self.catalog.create_table(&schema).await?; - } + // An existing table keeps its Iceberg schema — the columns and + // field ids the stream has evolved it to. `schema`'s columns come + // from discovery, which numbers them by position: after a column + // is dropped, the ones behind it would take others' field ids and + // land under the wrong Iceberg column. Nor does discovery's view + // of the source's *current* columns apply yet: the stream may + // still replay transactions from before a change. Its Relation + // messages evolve the schema in order ([`Self::apply_relation`]). + let schema = match self.catalog.load_table(&ident).await? { + None => { + self.catalog.create_table(&schema).await?; + schema + } + Some(meta) => TableSchema { + columns: meta.schema.columns, + ..schema + }, + }; // Notify the blob store that the table now exists in the // catalog. Static-creds backends ignore this; the // vended-credentials router uses it to load per-table STS @@ -873,6 +892,7 @@ impl Materializer { file_index, writer, gated_until_snapshot: false, + dropped: BTreeSet::new(), }, ); Ok(()) @@ -911,28 +931,28 @@ impl Materializer { } } - /// Apply a pgoutput Relation message: diff `incoming_columns` - /// against the registered schema for `ident`, build a - /// `Vec`, and call `Catalog::evolve_schema`. The - /// in-memory `TableEntry::schema` and `TableWriter` are - /// rebuilt from the post-evolution schema so subsequent - /// materialize cycles encode rows with the new shape. + /// Apply a pgoutput Relation message: bring the table's Iceberg + /// schema in line with `incoming_columns` ([`reconcile_columns`]), + /// call `Catalog::evolve_schema`, and rebuild the in-memory + /// `TableEntry::schema` and `TableWriter` from the result, so + /// subsequent materialize cycles encode rows with the new shape. /// - /// Detects three kinds of evolution: + /// - **AddColumn** for a new name, with a fresh field id. + /// - **DropColumn** (soft) for a non-PK column `incoming` lacks: it + /// stays, nullable, so older data files keep resolving. + /// - **PromoteColumnType** for a legal Iceberg promotion (int→long, + /// float→double, decimal precision increase). Illegal type + /// changes (e.g. long→int, text→int) are rejected with + /// `MaterializerError::Catalog` so the lifecycle fails loudly + /// rather than silently truncate or coerce downstream readers. + /// - **RenameColumn + AddColumn** for a column the source dropped + /// and re-added: Postgres gives it no values, so it's a new + /// column, and the dropped one keeps its values under another + /// name. /// - /// - **AddColumn**: a name in `incoming` that's not in our - /// schema. Field id is allocated by `apply_schema_changes` - /// from the schema's current high-water mark. - /// - **DropColumn**: a non-PK name in our schema that's not in - /// `incoming`. Soft-drop — column stays as nullable so older - /// data files keep resolving. - /// - **PromoteColumnType**: a name present in both, but - /// `incoming.ty != current.ty` *and* the change is a legal - /// Iceberg promotion (int→long, float→double, decimal - /// precision increase). Illegal type changes (e.g. long→int, - /// text→int) are rejected with `MaterializerError::Catalog` - /// so the lifecycle fails loudly rather than silently - /// truncate or coerce downstream readers. + /// The incoming `is_primary_key`/`nullable` are ignored: + /// pgoutput's key flag means "part of REPLICA IDENTITY" (every + /// column under `REPLICA IDENTITY FULL`), not "is primary key". /// /// **No-op when the table isn't registered** (the lifecycle /// only registers tables in YAML — incoming Relations for @@ -956,115 +976,26 @@ impl Materializer { Some(e) => e, None => return Ok(()), }; - - let mut current_names: std::collections::BTreeSet = entry - .schema - .columns - .iter() - .map(|c| c.name.clone()) - .collect(); - let incoming_names: std::collections::BTreeSet = - incoming_columns.iter().map(|c| c.name.clone()).collect(); - let current_by_name: std::collections::BTreeMap = entry - .schema - .columns + let incoming: Vec<(String, pg2iceberg_core::IcebergType)> = incoming_columns .iter() - .map(|c| (c.name.clone(), c)) + .map(|c| (c.name.clone(), c.ty)) .collect(); - - let mut changes: Vec = Vec::new(); - for c in incoming_columns { - match current_by_name.get(&c.name) { - None => { - // Force `nullable: true` for evolution-time adds. - // Two reasons: - // 1. Iceberg requires new columns to be optional - // so prior data files (which don't carry the - // column) read back as NULL. Pushing through a - // non-nullable add would break readers. - // 2. `c.nullable` is derived from pgoutput's - // Relation flag bit 0, which means "part of - // REPLICA IDENTITY" — *not* "is primary key". - // Under `REPLICA IDENTITY FULL` every column - // reports `is_replica_identity = 1`, so we'd - // incorrectly stamp every fresh column as - // non-nullable. The startup-time PK info from - // `discover_schemas` is the authoritative - // source; pgoutput's flag is wire-only signal. - changes.push(pg2iceberg_iceberg::SchemaChange::AddColumn { - name: c.name.clone(), - ty: c.ty, - nullable: true, - }); - } - Some(existing) => { - if existing.ty != c.ty { - // PK columns are part of the equality-delete - // predicate; promoting a PK type would - // invalidate every prior delete file's pk_key - // hash. Refuse — operators must re-snapshot - // for that case. - if existing.is_primary_key { - return Err(MaterializerError::Catalog(IcebergError::Other(format!( - "cannot promote primary-key column {}: {:?} → {:?} \ - (PK type is part of the equality-delete contract; \ - changing it requires a full re-snapshot)", - c.name, existing.ty, c.ty - )))); - } - if !pg2iceberg_iceberg::is_legal_type_promotion(existing.ty, c.ty) { - return Err(MaterializerError::Catalog(IcebergError::Other(format!( - "column {} type change {:?} → {:?} is not a legal Iceberg \ - promotion. Allowed: int→long, float→double, decimal \ - precision increase. Other changes (narrowing, cross-family) \ - require a full re-snapshot.", - c.name, existing.ty, c.ty - )))); - } - changes.push(pg2iceberg_iceberg::SchemaChange::PromoteColumnType { - name: c.name.clone(), - new_ty: c.ty, - }); - } - } - } - } - for c in &entry.schema.columns { - // PK columns are immutable in our model — pgoutput won't - // drop them anyway. Skip to avoid soft-dropping a PK - // (which would set `nullable = true` on the PK column). - if c.is_primary_key { - continue; - } - if !incoming_names.contains(&c.name) { - changes.push(pg2iceberg_iceberg::SchemaChange::DropColumn { - name: c.name.clone(), - }); - } - } - if changes.is_empty() { - return Ok(()); - } - - // Catalog-side first so a failure leaves the materializer's - // in-memory state unchanged (next Relation will re-trigger - // the diff). After success, rebuild the TableWriter so it - // encodes with the new column set. Use the Iceberg-side - // ident (`lookup_ident`) — `ident` is the PG-side key from - // the pgoutput Relation message, which the catalog - // wouldn't recognise. - self.catalog - .evolve_schema(&lookup_ident, changes.clone()) - .await?; - pg2iceberg_iceberg::apply_schema_changes(&mut entry.schema, &changes) + let changes = reconcile_columns(&entry.schema, &incoming, &entry.dropped) .map_err(MaterializerError::Catalog)?; - entry.writer = TableWriter::new(entry.schema.clone()); - // current_names is rebuilt next call from the updated schema; - // explicit insertion keeps Clippy happy + makes the - // post-state self-consistent within this call. - for c in &entry.schema.columns { - current_names.insert(c.name.clone()); + if !changes.is_empty() { + // Catalog-side first so a failure leaves the materializer's + // in-memory state unchanged (the next Relation re-triggers + // the diff). Use the Iceberg-side ident (`lookup_ident`) — + // `ident` is the PG-side key from the pgoutput Relation + // message, which the catalog wouldn't recognise. + self.catalog + .evolve_schema(&lookup_ident, changes.clone()) + .await?; + pg2iceberg_iceberg::apply_schema_changes(&mut entry.schema, &changes) + .map_err(MaterializerError::Catalog)?; + entry.writer = TableWriter::new(entry.schema.clone()); } + entry.dropped = absent_from(&entry.schema, &incoming); Ok(()) } @@ -1910,6 +1841,19 @@ fn expand_truncates( out } +/// The columns of `schema` that `source` doesn't have. +fn absent_from( + schema: &TableSchema, + source: &[(String, pg2iceberg_core::IcebergType)], +) -> BTreeSet { + schema + .columns + .iter() + .filter(|c| !source.iter().any(|(name, _)| *name == c.name)) + .map(|c| c.name.clone()) + .collect() +} + fn collect_toast_paths( rows: &[MaterializedRow], file_index: &FileIndex, diff --git a/crates/pg2iceberg-pg/src/prod/value_decode.rs b/crates/pg2iceberg-pg/src/prod/value_decode.rs index fce9211..4c0258a 100644 --- a/crates/pg2iceberg-pg/src/prod/value_decode.rs +++ b/crates/pg2iceberg-pg/src/prod/value_decode.rs @@ -39,11 +39,9 @@ pub fn decode_text(ty: PgType, bytes: &[u8]) -> Result { _ => return Err(bad()), }, PgType::Int2 => PgValue::Int2(raw.parse().map_err(|_| bad())?), - PgType::Int4 | PgType::Oid => { - // PG `oid` is unsigned 32-bit; reading as i64 then casting - // accepts wraparound for `oid` values above 2^31. - PgValue::Int4(raw.parse::().map_err(|_| bad())? as i32) - } + PgType::Int4 => PgValue::Int4(raw.parse().map_err(|_| bad())?), + // Unsigned 32-bit, so `long`: values past 2^31 don't fit `int`. + PgType::Oid => PgValue::Int8(raw.parse::().map_err(|_| bad())?.into()), PgType::Int8 => PgValue::Int8(raw.parse().map_err(|_| bad())?), PgType::Float4 => PgValue::Float4(raw.parse().map_err(|_| bad())?), PgType::Float8 => PgValue::Float8(raw.parse().map_err(|_| bad())?), diff --git a/crates/pg2iceberg-sim/src/postgres.rs b/crates/pg2iceberg-sim/src/postgres.rs index ce31ca0..53a099b 100644 --- a/crates/pg2iceberg-sim/src/postgres.rs +++ b/crates/pg2iceberg-sim/src/postgres.rs @@ -878,6 +878,22 @@ impl SimPostgres { } /// Snapshot of a table's rows in PK order, for tests / verify. + /// Whether `ident`'s relation changed (DDL, or an invalidation) after + /// its last row change. pgoutput tells a consumer about a schema + /// change only with the table's next change, so it hasn't told one. + pub fn relation_changed_since_last_change(&self, ident: &TableIdent) -> bool { + let s = self.state.lock().unwrap(); + let relation = s.wal.iter().rev().find_map(|e| match &e.kind { + WalKind::Relation { ident: i, .. } if i == ident => Some(e.lsn), + _ => None, + }); + let change = s.wal.iter().rev().find_map(|e| match &e.kind { + WalKind::Change(c) if &c.table == ident => Some(e.lsn), + _ => None, + }); + relation > change + } + pub fn read_table(&self, ident: &TableIdent) -> Result> { let s = self.state.lock().unwrap(); let t = s diff --git a/crates/pg2iceberg-tests/tests/dst.rs b/crates/pg2iceberg-tests/tests/dst.rs index 161a537..21a3bdd 100644 --- a/crates/pg2iceberg-tests/tests/dst.rs +++ b/crates/pg2iceberg-tests/tests/dst.rs @@ -1916,10 +1916,27 @@ impl Catalog for AuditedCatalog { } } +/// pgoutput tells pg2iceberg about a schema change only with the table's +/// next change; until a write comes, Iceberg can't reflect it — a column +/// dropped and re-added still shows its old values. That write comes +/// eventually: make it now, writing a row back unchanged, so the end +/// state compares what pg2iceberg can know. +fn write_after_schema_change(h: &DstHarness) { + if !h.db.relation_changed_since_last_change(&ident()) { + return; + } + if let Some(row) = h.db.read_table(&ident()).unwrap().into_iter().next() { + let mut tx = h.db.begin_tx(); + tx.update(&ident(), row); + tx.commit(Timestamp(0)).unwrap(); + } +} + fn check_invariants(h: &mut DstHarness) -> Result<(), String> { while h.backfilling { h.backfill_chunk(); } + write_after_schema_change(h); // Reach quiescence: drain WAL, flush, ack, then materialize until idle. // Loop because a flush may produce events the materializer hasn't seen. h.drive(); @@ -2209,6 +2226,7 @@ fn sort_by_pk(rows: &mut [Row]) { /// We keep invariants 1, 2, 3, and 5 — the headline correctness property /// (PG == Iceberg at quiescence) still holds, which is what matters. fn check_invariants_with_snapshot(h: &mut DstHarness) -> Result<(), String> { + write_after_schema_change(h); h.drive(); h.flush_and_ack(); for _ in 0..16 { @@ -2626,6 +2644,42 @@ fn re_added_column_does_not_bring_back_dropped_values() { check_invariants(&mut h).unwrap(); } +/// Rows staged before a column is dropped and re-added, but materialized +/// after: their values belong to the dropped column, not the new one. +/// (Fails today: schema changes apply when the Relation message arrives, +/// ahead of rows already staged, which are keyed by column name.) +#[test] +fn rows_staged_before_a_re_add_keep_their_values_out_of_the_new_column() { + let mut h = DstHarness::boot(); + h.run_step(&Step::Insert { id: 1, qty: 10 }); + h.run_step(&Step::DropNote); + h.run_step(&Step::AddNote); + h.run_step(&Step::Insert { id: 2, qty: 20 }); + check_invariants(&mut h).unwrap(); +} + +/// Compaction and TOAST resolution read older data files, which hold the +/// dropped column's values under the re-added column's name: they must +/// read columns by field id, or those values come back. +#[test] +fn reading_old_files_after_a_re_add_keeps_the_new_column_empty() { + let mut h = DstHarness::boot(); + h.run_step(&Step::Insert { id: 1, qty: 10 }); + h.run_step(&Step::Insert { id: 2, qty: 20 }); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + h.run_step(&Step::DropNote); + h.run_step(&Step::AddNote); + h.run_step(&Step::Insert { id: 3, qty: 30 }); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + h.run_step(&Step::ToastUpdate { id: 1, qty: 11 }); + h.run_step(&Step::DriveFlush); + h.run_step(&Step::MaterializerCycle); + h.run_step(&Step::Compact); + check_invariants(&mut h).unwrap(); +} + /// After a column ahead of others is dropped, a restarted materializer /// discovers the remaining columns at new positions; it must keep /// writing each value under its column's Iceberg field id.