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
1 change: 1 addition & 0 deletions crates/pg2iceberg-logical/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@

pub mod materializer;
pub mod pipeline;
pub mod relation_event;
pub mod runner;
pub mod sink;

Expand Down
148 changes: 75 additions & 73 deletions crates/pg2iceberg-logical/src/materializer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
//! 7. `Coordinator::set_cursor` to the highest end_offset processed.
//! 8. Update FileIndex with the new data file + removed PKs.

use crate::relation_event;
use async_trait::async_trait;
use bytes::Bytes;
use pg2iceberg_coord::{Coordinator, LogEntry};
Expand Down Expand Up @@ -461,15 +462,6 @@ pub struct Materializer<C: Catalog> {
/// controls how long a missed heartbeat survives before the
/// worker drops out of the active list.
distributed: Option<DistributedMode>,
/// PG → Iceberg ident translation, mirroring `Pipeline`'s map.
/// `apply_relation` is called with the PG-side ident (the
/// pgoutput stream's view), but `tables` is keyed by the
/// Iceberg-side ident the lifecycle registered with. Without
/// this translation, schema-evolution messages
/// (`ALTER TABLE ... ADD COLUMN`, type promotions, drops) would
/// silently no-op, and post-ALTER inserts would write to the
/// pre-ALTER iceberg schema.
table_translation: BTreeMap<TableIdent, TableIdent>,
}

/// Distributed-mode parameters. Built by
Expand Down Expand Up @@ -589,19 +581,6 @@ impl<C: Catalog> Materializer<C> {
meta_marker: None,
meta_recorder: None,
distributed: None,
table_translation: BTreeMap::new(),
}
}

/// Register a translation from the PG-side ident (carried in
/// pgoutput Relation messages) to the Iceberg-side ident this
/// materializer's `tables` map is keyed by. Lets schema-
/// evolution events from the WAL stream find their target
/// table even when `sink.namespace` differs from the PG schema.
/// No-op when the two idents are equal.
pub fn register_table_translation(&mut self, pg_ident: TableIdent, iceberg_ident: TableIdent) {
if pg_ident != iceberg_ident {
self.table_translation.insert(pg_ident, iceberg_ident);
}
}

Expand Down Expand Up @@ -843,7 +822,7 @@ impl<C: Catalog> Materializer<C> {
// 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`]).
// messages, staged in the log, evolve the schema in order.
let schema = match self.catalog.load_table(&ident).await? {
None => {
self.catalog.create_table(&schema).await?;
Expand Down Expand Up @@ -931,71 +910,81 @@ impl<C: Catalog> Materializer<C> {
}
}

/// 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.
/// Apply a staged relation event: bring the table's Iceberg schema
/// in line with the source columns it carries ([`reconcile_columns`]),
/// and rebuild the in-memory `TableEntry::schema` and `TableWriter`
/// from the result, so the rows after it encode with the new shape.
///
/// - **AddColumn** for a new name, with a fresh field id.
/// - **DropColumn** (soft) for a non-PK column `incoming` lacks: it
/// - **DropColumn** (soft) for a non-PK column the source 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.
/// `MaterializerError::Catalog` so the table 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.
///
/// 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
/// untracked tables, e.g. `_pg2iceberg.markers`, are silently
/// skipped).
///
/// **No-op when columns match.** pgoutput re-emits Relation
/// messages liberally (e.g. on every cache invalidation, even
/// for unchanged schemas); the diff just produces an empty
/// change list and we return without touching the catalog.
pub async fn apply_relation(
/// Reconciles against the catalog's schema, not the in-memory one:
/// the event may already have been applied — before a crash, whose
/// retry re-reads it, or by another worker — and applying it again
/// is then a no-op.
async fn apply_columns(
&mut self,
ident: &TableIdent,
incoming_columns: &[pg2iceberg_pg::RelationColumn],
columns: &relation_event::Columns,
) -> Result<()> {
// Pgoutput emits Relation messages keyed by PG schema +
// table name; our `tables` map is keyed by Iceberg-side
// ident. Translate when registered, otherwise fall through.
let lookup_ident = self.table_translation.get(ident).unwrap_or(ident).clone();
let entry = match self.tables.get_mut(&lookup_ident) {
Some(e) => e,
None => return Ok(()),
};
let incoming: Vec<(String, pg2iceberg_core::IcebergType)> = incoming_columns
.iter()
.map(|c| (c.name.clone(), c.ty))
.collect();
let changes = reconcile_columns(&entry.schema, &incoming, &entry.dropped)
let current = self
.catalog
.load_table(ident)
.await?
.ok_or_else(|| MaterializerError::UnknownTable(ident.clone()))?
.schema;
let entry = self
.tables
.get_mut(ident)
.ok_or_else(|| MaterializerError::UnknownTable(ident.clone()))?;
let changes = reconcile_columns(&current, columns, &entry.dropped)
.map_err(MaterializerError::Catalog)?;
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());
let current = if changes.is_empty() {
current
} else {
self.catalog.evolve_schema(ident, changes).await?.schema
};
entry.schema.columns = current.columns;
entry.writer = TableWriter::new(entry.schema.clone());
entry.dropped = absent_from(&entry.schema, columns);
Ok(())
}

/// Buffer a log entry's events into `unit`, applying its relation
/// events in place: rows before one are staged under the old schema
/// first.
async fn buffer_entry(
&mut self,
ident: &TableIdent,
unit: &mut Unit,
events: Vec<MatEvent>,
) -> Result<()> {
for evt in events {
if evt.op != Op::Relation {
unit.buf.push(evt);
continue;
}
if !unit.buf.is_empty() {
self.prepare_step(ident, unit).await?;
}
let columns = relation_event::decode(&evt.row).ok_or_else(|| {
MaterializerError::Blob(StreamError::Decode(format!(
"malformed relation event for {ident}: {:?}",
evt.row
)))
})?;
self.apply_columns(ident, &columns).await?;
}
entry.dropped = absent_from(&entry.schema, &incoming);
Ok(())
}

Expand Down Expand Up @@ -1404,6 +1393,19 @@ impl<C: Catalog> Materializer<C> {
for e in entries {
let bytes = self.blob_store.get(&e.s3_path).await?;
let events = decode_chunk(&bytes)?;
// A schema change starts its entry (the pipeline stages it
// so). Commit the rows before it — written under the old
// schema — before applying it: a crash in between then
// re-reads this entry, and applying the change again is a
// no-op. Within a transaction (DDL mid-transaction) the
// rows so far can only be staged as a step.
if events.first().is_some_and(|e| e.op == Op::Relation) {
if is_tx_boundary(unit.buf.last(), events.first()) {
folded += self.commit_unit(ident, &mut unit).await?;
} else {
self.prepare_step(ident, &mut unit).await?;
}
}
if !unit.buf.is_empty() && unit.buf.len() + events.len() > self.batch_rows {
if is_tx_boundary(unit.buf.last(), events.first()) {
folded += self.commit_unit(ident, &mut unit).await?;
Expand All @@ -1418,7 +1420,7 @@ impl<C: Catalog> Materializer<C> {
}
}
events_read += events.len();
unit.buf.extend(events);
self.buffer_entry(ident, &mut unit, events).await?;
unit.end_offset = Some(e.end_offset);
after = e.end_offset;
}
Expand Down
71 changes: 60 additions & 11 deletions crates/pg2iceberg-logical/src/pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,16 @@
//! `flush()` drains, uploads, claims, and advances `flushedLSN` — the last
//! step gated by [`CoordCommitReceipt`].

use crate::relation_event;
use crate::sink::{FlushOutput, Sink, SinkError, TableChunk};
use async_trait::async_trait;
use pg2iceberg_coord::{
CommitBatch, CoordCommitReceipt, CoordError, Coordinator, MarkerInfo, OffsetClaim,
};
use pg2iceberg_core::metrics::{names, Labels};
use pg2iceberg_core::{ColumnName, Lsn, Metrics, NoopMetrics, Op, PgValue, TableIdent};
use pg2iceberg_core::{
ChangeEvent, ColumnName, Lsn, Metrics, NoopMetrics, Op, PgValue, TableIdent, Timestamp,
};
use pg2iceberg_pg::DecodedMessage;
use pg2iceberg_stream::codec::EncodedChunk;
use pg2iceberg_stream::{BlobStore, StreamError};
Expand Down Expand Up @@ -149,6 +152,10 @@ pub struct Pipeline<C: Coordinator + ?Sized> {
/// This pipeline consumes the replication stream (see
/// [`Self::track_replication`]).
replication: bool,
/// The open transaction: xid and `Begin`'s LSN.
open_tx: Option<(u32, Lsn)>,
/// Each table's columns as last staged (see [`Self::stage_relation`]).
relations: BTreeMap<TableIdent, relation_event::Columns>,
}

impl<C: Coordinator + ?Sized> Pipeline<C> {
Expand Down Expand Up @@ -192,6 +199,8 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
primary_keys: BTreeMap::new(),
table_translation: BTreeMap::new(),
replication: false,
open_tx: None,
relations: BTreeMap::new(),
}
}

Expand Down Expand Up @@ -262,8 +271,12 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
return Ok(());
}
match msg {
DecodedMessage::Begin { xid, .. } => self.sink.begin_tx(xid),
DecodedMessage::Begin { xid, final_lsn } => {
self.open_tx = Some((xid, final_lsn));
self.sink.begin_tx(xid);
}
DecodedMessage::Commit { xid, commit_lsn } => {
self.open_tx = None;
// Drain any markers observed in this tx before
// committing. Flushed atomically with the rest of
// the tx via the next claim_offsets call.
Expand Down Expand Up @@ -386,10 +399,8 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
self.sink.record_change(evt)?;
self.spill_if_full(xid).await?;
}
DecodedMessage::Relation { .. } => {
// Schema evolution is applied via
// `Materializer::apply_relation`, which the lifecycle
// calls *before* this dispatch. Nothing to do here.
DecodedMessage::Relation { ident, columns } => {
self.stage_relation(ident, &columns)?;
}
DecodedMessage::Keepalive { wal_end, .. } => {
// Only trustworthy between transactions: pgoutput sends a
Expand All @@ -404,6 +415,47 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
Ok(())
}

/// Stage the table's columns, if they changed, as a relation event
/// where the stream put the Relation message — before the
/// transaction's changes to the table — so the materializer applies
/// the schema change between the rows staged before and after it.
/// Applied when the message arrived instead, ahead of rows already
/// staged under the old schema, a dropped and re-added column would
/// take those rows' values for the dropped one.
fn stage_relation(
&mut self,
ident: TableIdent,
columns: &[pg2iceberg_pg::RelationColumn],
) -> Result<()> {
let table = self.table_translation.get(&ident).cloned().unwrap_or(ident);
if self.markers_table.as_ref() == Some(&table) {
return Ok(());
}
let columns: relation_event::Columns =
columns.iter().map(|c| (c.name.clone(), c.ty)).collect();
// pgoutput resends a table's Relation, unchanged, at every new
// session and cache invalidation; those say nothing new.
if self.relations.get(&table) == Some(&columns) {
return Ok(());
}
let (xid, lsn) = match self.open_tx {
Some((xid, lsn)) => (Some(xid), lsn),
None => (None, Lsn::ZERO),
};
self.sink.record_change(ChangeEvent {
table: table.clone(),
op: Op::Relation,
lsn,
commit_ts: Timestamp(0),
xid,
before: None,
after: Some(relation_event::encode(&columns)),
unchanged_cols: Vec::new(),
})?;
self.relations.insert(table, columns);
Ok(())
}

/// Drain all committed-but-unflushed transactions: encode → upload →
/// `claim_offsets` → advance `flushedLSN` via the receipt. No-op if
/// nothing is ready (no staged events, no markers awaiting emission,
Expand Down Expand Up @@ -520,11 +572,8 @@ impl<C: Coordinator + ?Sized> Pipeline<C> {
self.blob_store.put(&path, chunk.bytes).await?;
let mut labels = Labels::new();
labels.insert("table".into(), table.name.clone());
self.metrics.counter(
names::PIPELINE_ROWS_STAGED_TOTAL,
&labels,
chunk.record_count,
);
self.metrics
.counter(names::PIPELINE_ROWS_STAGED_TOTAL, &labels, chunk.row_count);
Ok(OffsetClaim {
table,
record_count: chunk.record_count,
Expand Down
43 changes: 43 additions & 0 deletions crates/pg2iceberg-logical/src/relation_event.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
//! A source table's columns, as staged in the change log.
//!
//! The pipeline stages each pgoutput Relation message that changes a
//! table's columns as an `Op::Relation` event, in stream order, so the
//! materializer applies schema changes between the rows staged before
//! and after them. Its row carries the column list under one key, in the
//! source's column order (which re-add detection reads — see
//! [`pg2iceberg_iceberg::reconcile_columns`]).

use pg2iceberg_core::{ColumnName, IcebergType, PgValue, Row};

const COLUMNS: &str = "columns";

/// A table's columns: name and type, in the source's order.
pub type Columns = Vec<(String, IcebergType)>;

pub fn encode(columns: &Columns) -> Row {
let json = serde_json::to_string(columns).expect("columns serialize");
Row::from([(ColumnName(COLUMNS.into()), PgValue::Json(json))])
}

/// `None` if `row` isn't an encoded column list.
pub fn decode(row: &Row) -> Option<Columns> {
match row.get(&ColumnName(COLUMNS.into()))? {
PgValue::Json(json) => serde_json::from_str(json).ok(),
_ => None,
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn columns_round_trip_in_order() {
let cols: Columns = vec![
("id".into(), IcebergType::Int),
("qty".into(), IcebergType::Long),
("note".into(), IcebergType::String),
];
assert_eq!(decode(&encode(&cols)), Some(cols));
}
}
Loading
Loading