# Logical Replication

> Essential Postgres logical replication concepts for working with Supabase ETL.

- Canonical HTML: https://supabase.github.io/etl/explanation/concepts/
- Agent-readable Markdown: https://supabase.github.io/etl/explanation/concepts.md
- Source: https://github.com/supabase/etl/blob/main/site/content/docs/explanation/concepts.md

Read this first if Postgres logical replication is new to you.

## What is Logical Replication? [#what-is-logical-replication]

Postgres supports two types of replication:

| Type         | What it copies                               | Use case                         |
| ------------ | -------------------------------------------- | -------------------------------- |
| **Physical** | Exact byte-for-byte copy of data files       | Disaster recovery, read replicas |
| **Logical**  | Decoded row changes (INSERT, UPDATE, DELETE) | Data integration, ETL, CDC       |

ETL uses logical replication because decoded row changes can be sent to systems
other than Postgres.

## The Write-Ahead Log (WAL) [#the-write-ahead-log-wal]

Before Postgres modifies data on disk, it first writes the change to the &#x2A;*Write-Ahead Log (WAL)**. This guarantees durability: if Postgres crashes, it can replay the WAL to recover.

<Mermaid
  chart="flowchart LR
A[Transaction commits] --> B[Written to WAL] --> C[Later flushed to data files]"
/>

For logical replication, Postgres decodes the WAL back into **logical changes**:

<Mermaid
  chart="flowchart LR
A[WAL bytes] --> B[&#x22;Decoder (pgoutput)&#x22;] --> C[INSERT/UPDATE/DELETE events]"
/>

ETL receives these decoded events and forwards them to downstream consumers.

### WAL Level [#wal-level]

Postgres must be configured to record enough information for logical decoding:

```ini
# In postgresql.conf
wal_level = logical
```

With `wal_level = logical`, Postgres records additional metadata needed to reconstruct row changes. Lower levels (`replica`, `minimal`) **do not capture enough detail**.

## Publications [#publications]

A **publication** defines which tables to replicate. Think of it as a filter that says "replicate changes from these tables."

```sql
-- Replicate specific tables
CREATE PUBLICATION my_publication FOR TABLE users, orders;

-- Replicate all tables (use with caution)
CREATE PUBLICATION my_publication FOR ALL TABLES;
```

When you create an ETL pipeline, you specify which publication to consume.
&#x2A;*Only tables and operations selected by that publication are replicated.**

### What Publications Control [#what-publications-control]

* **Which tables**: Only tables in the publication are replicated
* **Which operations**: You can filter to only INSERT, UPDATE, or DELETE
* **Which columns** (Postgres 15+): Replicate only specific columns
* **Which rows** (Postgres 15+): Filter rows with a WHERE clause

## Replication Slots [#replication-slots]

A **replication slot** is a bookmark that tracks how far a consumer has read in the WAL.

### Why Slots Exist [#why-slots-exist]

Without slots, Postgres would delete old WAL files when it no longer needs them for crash recovery. If ETL disconnects temporarily, it needs those WAL files to catch up when it reconnects.

Replication slots tell Postgres: &#x2A;*"Don't delete WAL files until this consumer has processed them."**

```sql
-- View existing slots
SELECT slot_name, confirmed_flush_lsn, active
FROM pg_replication_slots;
```

### How ETL Uses Slots [#how-etl-uses-slots]

ETL creates replication slots automatically:

| Slot                                               | Purpose                           |
| -------------------------------------------------- | --------------------------------- |
| `supabase_etl_apply_{pipeline_id}`                 | Main slot for ongoing replication |
| `supabase_etl_table_sync_{pipeline_id}_{table_id}` | Temporary slots for initial sync  |

The Apply Worker uses one persistent slot. Table Sync Workers create temporary slots during initial sync, then delete them.

### Slot Risks [#slot-risks]

Slots prevent WAL cleanup. If ETL stops consuming because of crashes, network issues, or a slow consumer, WAL files accumulate on disk. &#x2A;*This can fill your disk.**

To mitigate this risk:

* Monitor slot lag with `pg_replication_slots`
* Set `max_slot_wal_keep_size` to limit WAL retention
* Alert when slots fall behind

See [Configure Postgres](https://supabase.github.io/etl/guides/configure-postgres.md#wal-buildup-and-disk-usage) for details.

## The pgoutput Decoder [#the-pgoutput-decoder]

When Postgres decodes WAL for logical replication, it uses a **decoder plugin**. ETL uses `pgoutput`, Postgres's built-in decoder.

The decoder transforms binary WAL records into structured messages:

| Message    | Meaning                       |
| ---------- | ----------------------------- |
| `BEGIN`    | Transaction started           |
| `RELATION` | Table schema (columns, types) |
| `INSERT`   | Row added                     |
| `UPDATE`   | Row modified                  |
| `DELETE`   | Row removed                   |
| `TRUNCATE` | Table cleared                 |
| `COMMIT`   | Transaction completed         |

ETL receives these messages and converts them to events.

## Why Two Phases? [#why-two-phases]

ETL replicates data in two phases: **initial sync** and **ongoing
replication**.

### Phase 1: Initial Sync [#phase-1-initial-sync]

Logical replication only captures **changes**. It does not know about data that existed before replication started.

So ETL first copies all existing rows using Postgres's `COPY` command:

1. Create replication slot (captures consistent snapshot point)
2. COPY all rows from the table
3. Begin ongoing replication from the snapshot point

The slot ensures **no changes are lost** between the snapshot and ongoing
replication.

### Phase 2: Ongoing Replication [#phase-2-ongoing-replication]

After initial sync, ETL begins ongoing replication. It captures subsequent
changes from the WAL and delivers them to the destination:

<Mermaid
  chart="flowchart LR
A[Postgres WAL] --> B[Decoder] --> C[ETL] --> D[Destination]"
/>

Each change is delivered as an `Event` through `write_events()`.

Large tables can spend significant time in initial sync. ETL exposes separate
`write_table_rows()` and `write_events()` methods so destinations can optimize
initial-copy rows and change events independently.

## Replica Identity [#replica-identity]

**REPLICA IDENTITY** controls what data Postgres includes in UPDATE and DELETE events.

### The Problem [#the-problem]

When a row is updated or deleted, downstream systems need enough old-row information to identify which source row changed.

The important nuance is that PostgreSQL does **not** always send an old-side tuple for `UPDATE`.
Under key-based replica identity, it only sends a key image when it is needed.
For `DELETE`, PostgreSQL sends an old-side tuple whenever the delete is publishable.

This means replica identity is both a **PostgreSQL logging rule** and a **consumer contract** for
downstream consumers. It determines whether an event contains enough old-row
data to match an existing row, detect key changes, or compare before-and-after
values.

### Settings [#settings]

```sql
-- See current setting (d=default, f=full, n=nothing, i=index)
SELECT relname, relreplident FROM pg_class WHERE relname = 'your_table';

-- Change setting
ALTER TABLE your_table REPLICA IDENTITY FULL;
```

| Setting                         | Published `UPDATE` payload                                                                                            | Published `DELETE` payload                                   | Notes                                                             |
| ------------------------------- | --------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------ | ----------------------------------------------------------------- |
| `DEFAULT` with a primary key    | Old primary-key columns only when PostgreSQL determines the old key must be logged; otherwise no old tuple            | Old primary-key columns                                      | Most tables with a primary key                                    |
| `DEFAULT` without a primary key | Source `UPDATE` is rejected when the table publishes updates                                                          | Source `DELETE` is rejected when the table publishes deletes | Equivalent to having no usable replica identity for update/delete |
| `FULL`                          | Full old row                                                                                                          | Full old row                                                 | Use when consumers need full old-row images                       |
| `NOTHING`                       | Source `UPDATE` is rejected when the table publishes updates                                                          | Source `DELETE` is rejected when the table publishes deletes | Suitable only when updates/deletes are not published              |
| `USING INDEX`                   | Old replica-identity index columns only when PostgreSQL determines the old key must be logged; otherwise no old tuple | Old replica-identity index columns                           | Tables whose replication identity differs from the primary key    |

### Impact on ETL [#impact-on-etl]

ETL preserves PostgreSQL's old-row semantics in update and delete events:

```rust
pub old_table_row: Option<OldTableRow>
```

* `Some(OldTableRow::Key(row))` means PostgreSQL sent only the replica-identity columns, normalized into replicated table-column order.
* `Some(OldTableRow::Full(row))` means PostgreSQL sent the full old row.
* `None` means PostgreSQL did not send an old-side tuple for that update. This
  is normal under `DEFAULT` or `USING INDEX` when PostgreSQL determines no
  old-side image is required.

For `FULL`, PostgreSQL sends a full old row for every published update and delete.
For `DELETE`, valid pgoutput messages always include either a full old row or a
key image. `REPLICA IDENTITY NOTHING`, and `DEFAULT` on a table without a primary
key, do not produce update/delete events when those actions are published; the
source statement is rejected instead.
The Rust event API keeps the old-row fields optional at the boundary, but those
`None` cases are broader than the PostgreSQL pgoutput shapes described here.

**TOAST adds one more wrinkle.** PostgreSQL can mark unchanged toasted update values
as `UnchangedToast` instead of resending the value. ETL can reconstruct those
values only if the old-side row image contains them, so tables with toasted
columns can produce partial update rows unless they use `REPLICA IDENTITY FULL`
or the missing values are present in a logged key image.

If you need **old values** for auditing, comparison, complete replacement rows,
or reliable reconstruction of unchanged toasted columns, set
`REPLICA IDENTITY FULL` on those tables. If a consumer only needs stable key
values, `DEFAULT` with a primary key or `USING INDEX` is usually enough, but
update events will not always include `old_table_row`.

## LSN (Log Sequence Number) [#lsn-log-sequence-number]

Every position in the WAL has a unique **LSN** - a monotonically increasing pointer.

```text
Format: 0/16B3748 (segment/offset)
```

### LSNs in Events [#lsns-in-events]

Sequenced ETL events include a commit LSN and transaction-local ordinal:

| Field        | Meaning                                       |
| ------------ | --------------------------------------------- |
| `commit_lsn` | LSN of the commit message in the WAL          |
| `tx_ordinal` | Zero-based event order within the transaction |

Multiple events in the same transaction share the same `commit_lsn`; their
`tx_ordinal` values distinguish their order. Relation events are connection-local
metadata and do not have an event sequence key.

## Persisted State [#persisted-state]

ETL persists the state needed to resume safely after a restart:

ETL stores:

| State                            | Purpose                                                                      |
| -------------------------------- | ---------------------------------------------------------------------------- |
| Table state                      | Track each table from initial sync through ongoing replication               |
| Persisted replication checkpoint | Resume workers from a safe replay frontier                                   |
| Table schemas                    | Decode events against the correct versioned schema                           |
| Destination table metadata       | Track destination table IDs, applied schema snapshots, and replication masks |

The built-in `PostgresStore` persists to your Postgres database and runs its
state-store migrations when it is created. If the pipeline reads from a
read-only replica, configure `store_pg_connection` to point at a writable
Postgres endpoint for this state. `MemoryStore` is for testing only - state is
lost on restart. `Pipeline::start()` runs the ETL source migrations that install
schema helpers and the DDL event trigger before replication begins. See
[Architecture](https://supabase.github.io/etl/explanation/architecture.md) for the worker lifecycle and
[Configure Postgres](https://supabase.github.io/etl/guides/configure-postgres.md) for production settings.