Skip to content
Merged
Show file tree
Hide file tree
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
18 changes: 9 additions & 9 deletions crates/pg2iceberg-core/src/typemap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,8 +87,9 @@ pub const DEFAULT_NUMERIC_SCALE: u8 = 18;
pub fn map_pg_to_iceberg(pg: PgType) -> Result<Mapped, MapError> {
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),
Expand Down Expand Up @@ -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]
Expand Down
278 changes: 278 additions & 0 deletions crates/pg2iceberg-iceberg/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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<String>,
) -> Result<Vec<SchemaChange>> {
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<usize> = 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<()>;
Expand Down Expand Up @@ -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"));
}
}
Loading
Loading