Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

CDC Reference

Rivet’s change-data-capture (CDC) reads a source’s transaction log — not the tables — and emits each INSERT / UPDATE / DELETE as a row change. Because it tails the log that the database already writes for replication and durability, it adds almost no load to the OLTP path: no table scan, no locks, no read snapshot (see Why CDC is gentle on the source).

Status. All three SQL engines support NDJSON streaming and typed Parquet/CSV --output (real Timestamp / Date32 / Decimal128 columns). This page documents all three so the permissions are clear up front — they are the part operators most need to get right.

MongoDB also has CDC — via change streams, with a different setup (a replica set, not per-table grants) and the JSON-blob document image rather than typed columns. It has its own reference: mongodb.md.

The command

# stream changes as NDJSON to stdout (no schema resolution, fewest privileges)
rivet cdc --source 'mysql://rivet_cdc:***@127.0.0.1:3306/app' --table orders

# write typed Parquet files (one row per change, after-image / upsert shape)
rivet cdc --source 'mysql://rivet_cdc:***@127.0.0.1:3306/app' \
          --table orders --output ./cdc-out --format parquet \
          --checkpoint ./orders.ckpt --rollover 100000

Prefer --source-env VAR or --source-file path over an inline URL outside local dev — the URL is otherwise visible in ps / shell history.

The rivet cdc CLI is loopback-only: it carries no TLS configuration, so the TLS gate refuses any remote (non-loopback) host before connecting. For a remote source, use the config-driven rivet run path with a source.tls: block (see From config).

flagmeaning
--server-idreplica id for the binlog connection (MySQL). Must be unique — distinct from the source and every real replica. Default 4271.
--checkpoint PATHpersist/resume the log position. Omission semantics differ per engine: MySQL and MongoDB tail from the current position without checkpointing; PostgreSQL anchors server-side at the slot regardless (omission merely disarms the slot-loss hard error the checkpoint enables); SQL Server with no checkpoint re-reads the entire retained change table on every run. Keep it set. On the first checkpointed MySQL run the open position is persisted immediately (the client-side analogue of PostgreSQL’s slot pinning at creation), so an idle first run still anchors the resume position — without it, changes landing between two idle scheduler cycles would be skipped.
--table NAMEonly emit this table (repeatable for NDJSON; exactly one required for --output, whose schema is resolved from the source).
--output DIRwrite typed Parquet/CSV files instead of NDJSON.
--max-events Nstop after N changes; without it the default bounded run drains to the log end as of open and exits (stream until interrupted only with --stream). The checkpoint is saved at transaction-commit boundaries (never mid-transaction), so an interrupted run resumes from the last fully committed transaction — re-reading, never skipping, a partially processed one.
--rollover Nrows per output part file (default 100000); also rolls at a transaction boundary, never splitting one. This is the file-size ⇄ memory dial: larger ⇒ fewer, bigger files but more drain memory (the PostgreSQL peek reads a part’s worth per batch, so drain RSS is O(rollover) — ≈28 MB + 1.3 KB × rollover). Raise it to cut file count on a big host; lower it to cap memory on a small extractor. (Config: cdc.rollover.)
--slot NAMEPostgreSQL logical slot (default rivet_slot; created if absent).
--capture-instance NAMESQL Server CDC capture instance (e.g. dbo_orders) — required for sqlserver://.
--streamOpt out of the default bounded run and stream continuously (a long-lived daemon). By default rivet cdc catches up to the source’s log end as of the moment the run opened, then exits instead of streaming — this is the scheduler-friendly model, so no flag is needed for it. Every engine pins that boundary at open (PostgreSQL: pg_current_wal_lsn(); MySQL: the binlog coordinates, plus BINLOG_DUMP_NON_BLOCK as the catch-up backstop; SQL Server: fn_cdc_get_max_lsn(); MongoDB: the cluster operationTime; Oracle: the current SCN), so a hot table whose writers outpace the drain cannot keep the run alive chasing a moving log end — the run’s work is O(backlog at open), and everything committed after the boundary is picked up by the next run from the checkpoint. With --max-events N, the bounded run stops at the smaller of “N events” or the boundary — so it never blocks waiting for the N-th event. Passing --stream removes the open-time boundary, but what that means is engine-specific: MySQL genuinely stays up (the binlog dump blocks on an idle source), and so does MongoDB (the change stream blocks awaiting events; it ends only if the stream is invalidated or closed); PostgreSQL and SQL Server are poll adapters that still exit on catch-up — one unbounded pass, not a daemon. Oracle refuses --stream (and cdc.until_current: false): LogMiner here is always a bounded drain to the SCN current at open. On PostgreSQL and SQL Server a continuous pipeline needs an external supervisor re-running the command (or just the default bounded model on a schedule). On a PostgreSQL standby (PG 16+ logical decoding) the ceiling query (pg_current_wal_lsn()) is unavailable during recovery, so the default bounded run fails loudly at open — pass --stream, or point the source at the primary.

The engine is chosen from the URL scheme (mysql:// / postgresql:// / sqlserver:// / mongodb://) by create_change_stream, the CDC sibling of the batch create_source. With --output, each part goes through the same commit path the batch export uses (ADR-0004) and a manifest.json + _SUCCESS is written at clean end — but the CLI’s --output is a local directory only (it is wired to the local destination; a gs://…/s3://… string would be taken as a literal local path). For a cloud destination, use the config path (mode: cdc with a destination: block) below. Typed columns (real Timestamp / Date32 / Decimal128, not strings) flow through RivetValue structural typing — for all three engines (MySQL binlog values, PostgreSQL test_decoding parse, SQL Server change-table ColumnData).

From config (rivet run)

CDC also runs as an export in a config, so a scheduled rivet run captures changes alongside batch exports and records the run the same way:

source:
  type: mysql
  url_env: DATABASE_URL          # credentials out of the file
  tls: { mode: verify-full }     # required for a remote host (see below)
exports:
  - name: orders_cdc
    table: orders
    mode: cdc
    format: parquet
    cdc:
      checkpoint: /var/lib/rivet/orders.ckpt
      until_current: true        # drain to now and exit — for a scheduler
      # per-engine, all optional:
      server_id: 4271            # MySQL replica id
      slot: rivet_orders         # PostgreSQL logical slot
      capture_instance: dbo_orders  # SQL Server (required for sqlserver://)
    destination: { type: gcs, bucket: my-bucket, prefix: cdc/orders }
rivet run --config cdc.yaml      # captures, writes typed Parquet, records the run
rivet metrics -c cdc.yaml        # the CDC run appears with mode=cdc, like a batch

A mode: cdc export reuses the export’s table, destination, and format; the cdc: block carries only the CDC-specific knobs.

initial: snapshot — the safe switch, enforced by construction. On the first run (no anchor yet) rivet performs, in order: ① create the resume anchor (PostgreSQL slot / MySQL binlog checkpoint / SQL Server max-LSN checkpoint), ② run a full batch snapshot of each table into <destination>[/<table>]/snapshot/ (its own parts + manifest.json + _SUCCESS), ③ drain the change stream. Because the anchor predates the snapshot read, a change landing mid-snapshot appears in both the snapshot and the stream — an overlap the PK + __op dedupe absorbs, never a gap. Subsequent runs skip the snapshot because the state DB records it as done (cdc_snapshot — the authoritative signal; the snapshot/_SUCCESS marker remains a legacy co-signal, so re-snapshotting requires clearing BOTH) and go straight to draining; a run that crashes mid-snapshot re-snapshots on retry (the anchor stays put, so nothing is lost). Once any snapshot completed, a MISSING server-side anchor (a dropped PostgreSQL slot) is a loud error, never a silent re-anchor — see “A vanished slot” below. Load order downstream: the snapshot prefix as the base table, then MERGE the CDC parts. MySQL / SQL Server require cdc.checkpoint: with initial: snapshot (it is the anchor); PostgreSQL anchors in the slot.

  - name: orders_cdc
    table: orders
    mode: cdc
    format: parquet
    cdc: { initial: snapshot, checkpoint: /var/lib/rivet/orders.ckpt, until_current: true }
    destination: { type: gcs, bucket: my-bucket, prefix: cdc/orders }

cdc.backfill: — the baseline by reference, for tables that disagree. initial: snapshot synthesizes one single-connection mode: full scan per captured table. That is right for a small table and wrong for a large one — a 313M-row table read end to end on one statement runs into tuning.statement_timeout_s (300s under the balanced profile) long before it finishes, while the same table as a batch export, keyset-paged with parallel: 4, takes ~22 minutes. And a multiplex stream’s tables do not agree on how they are read: one has a unique id and keysets, another has only a non-unique index and must range-chunk.

So the baseline is declared by REFERENCE — each captured table names the ordinary batch export that already describes how to read it:

exports:
  - name: orders                 # an ordinary export; `rivet init` already writes it
    table: orders
    mode: chunked
    chunk_by_key: id             # keyset
    parallel: 4
    chunk_checkpoint: true       # the baseline is resumable
    format: parquet
    destination: { type: gcs, bucket: my-bucket, prefix: exports/orders/ }
    columns: { price: decimal(10,2) }

  - name: app_cdc
    tables: [orders]
    mode: cdc
    format: parquet
    cdc:
      checkpoint: /var/lib/rivet/app.ckpt
      backfill: auto             # or: [orders, …]
    destination: { type: gcs, bucket: my-bucket, prefix: cdc/ }

auto pairs each entry of tables: with the export whose table: names it; a list names them explicitly. One rivet run then does anchor → baseline → drain, and the ordering is what makes it safe: the anchor is taken before the first row is read, so a row changed mid-baseline also arrives on the stream and the current-state view keeps the higher (__pos, __seq).

What the leg borrows and what stays its own is the whole design. Borrowed: mode, chunk_by_key / chunk_column, page size, parallel, chunk_checkpoint, tuning and the column types. Its own: the name, the <destination>/<table>/snapshot/ prefix, the format and the meta columns — so the referenced export contributes a recipe, never a second load target, and every load invariant that holds for a synthesized leg holds unchanged here. Types MERGE (the recipe’s, with a qualified "table.column" key on the CDC export still winning); a column both sides declare differently is refused, because the two legs write into one <table>__changes and one column cannot have two types.

The run loop skips an export that is named as a backfill recipe, so a full rivet run reads each table once — rivet run -e orders still exports it on its own. An interrupted baseline resumes on the next plain rivet run from its chunk checkpoints (range-chunked and keyset legs alike, both live-proven against a crash after the first page — no --resume, no synthesized name; a leg whose recipe changed after the crash names rivet state reset-chunks -e <leg>) and leaves the anchor alone; once a table’s baseline is recorded (per table, in the state DB), later runs go straight to the drain. cdc.initial: and cdc.backfill: both describe the first run’s baseline, so config load refuses the pair.

Multiple CDC exports: each owns its stream resources. A PostgreSQL slot has ONE consumer (a shared slot is advanced past changes the other export never read), a MySQL server_id has ONE connection (the server kills the older one), and a checkpoint file has ONE writer. Config validation rejects two mode: cdc exports that resolve to the same slot / server_id / checkpoint — including the defaults (rivet_slot, 4271): a multi-table CDC config must set them explicitly per export:

exports:
  - name: orders_cdc
    table: orders
    mode: cdc
    cdc: { slot: rivet_orders, checkpoint: /var/lib/rivet/orders.ckpt }
    ...
  - name: users_cdc
    table: users
    mode: cdc
    cdc: { slot: rivet_users, checkpoint: /var/lib/rivet/users.ckpt }
    ...

Or multiplex: several tables through ONE stream (tables:). N single-table exports cost N slots — and PostgreSQL decodes the WAL once per slot (MySQL: N binlog connections). One export with tables: rides a single slot/connection and a single checkpoint, and routes each table’s changes to its own sub-prefix (<destination>/<table>/, each with its own parts + manifest.json + _SUCCESS — exactly like N exports, minus the N−1 slots):

exports:
  - name: app_cdc
    tables: [orders, users, payments]
    mode: cdc
    format: parquet
    cdc: { slot: rivet_app, checkpoint: /var/lib/rivet/app.ckpt, until_current: true }
    destination: { type: local, path: /data/cdc }   # → /data/cdc/orders/, /data/cdc/users/, …

The resume position is a property of the stream, so the at-least-once sequence generalises: every table’s buffered part is flushed before the one checkpoint/ack advances — a crash mid-roll re-reads for all tables rather than losing any one of them.

columns: overrides on a multi-table export support two key shapes: a bare column name applies to every captured table that has it, and a qualified "table.column" key targets one table and wins over the bare key there — so same-named columns needing different treatments never collide:

    columns:
      amount: "decimal(20,4)"          # every table's `amount`
      "legacy_orders.amount": text     # …except this one

A qualified key naming a table the export does not capture is a config error (a typo must fail at load, never silently miss its target).

Whole-database CDC across engines — and why the config shape differs. tables: multiplexing is PostgreSQL/MySQL only, and that is a property of the engine, not a rivet or driver limit. MySQL exposes one server-wide binlog and PostgreSQL one logical slot — a single stream carrying every table’s changes, so rivet reads it once and routes by table (and two exports sharing the slot/server_id would collide — exactly the scarce resource the one-stream form conserves). SQL Server has no such stream: CDC is enabled per table (sys.sp_cdc_enable_table) and read through a per-capture-instance function (cdc.fn_cdc_get_all_changes_<instance>) — there is no server-wide “all changes” surface to tap, so capture is inherently per-table. Use one export per table there, each with its own capture_instance; sharing a capture_instance between exports is safe (the change-table poll is read-only and resume state lives in the per-export checkpoint), and per-table exports never collide (no slot / server_id — the change tables are populated by one shared capture Agent regardless of how many readers).

The operator flow is identical across the three SQL engines, only the config shape differs (MongoDB’s whole-database change stream uses the same mode: cdc config shape — see mongodb.md):

  • rivet init --mode cdc (no --table) scaffolds the whole database on every engine — one tables: export on MySQL (and on PostgreSQL when every table is in the public schema; mixed schemas fall back to per-table exports), one export per table (distinct capture_instance) on SQL Server — so you never hand-list tables. Over two or more tables on MySQL/PostgreSQL the stream gets backfill: auto and one batch recipe per table (the baseline read), and its load: may carry tables: { <table>: { pk, partition, cluster_by, … } } so each captured table has its own warehouse shape.
  • rivet run -c <config> drains the whole set; add --parallel-export-processes to run SQL Server’s per-table exports concurrently.
  • rivet validate -c <config> descends into every table’s prefix and its initial: snapshot sub-dataset on every engine, so one command certifies the whole stream.

The per-table vs one-stream split is connection/resource topology, not throughput or memory: drain RSS is O(part rollover) per stream on all three engines, independent of table count. Each run produces the standard per-export summary block and an export_metrics row (rows / files / bytes / duration / status), so CDC shows up in rivet metrics and the run aggregate exactly like a batch export. TLS: unlike the CLI (which is loopback-only), the config path passes source.tls to the change stream — so a remote source over TLS requires the tls: block, and a remote host without it is refused before any connection (the same gate the batch path uses).

The four models

Rivet normalises four different source mechanisms behind one ChangeStream:

enginemechanismmodel
MySQLbinlog (ROW) streamed as a replicapush — the client reads the log directly
PostgreSQLlogical replication slot (test_decoding)poll the slot via pg_logical_slot_peek_changes()
SQL Servercdc.* change tables the capture Agent extractspoll the change function by LSN window
MongoDBwhole-database change stream (db.watch() over the oplog)tailable stream; the resume token checkpoints the position (JSON-blob image — see mongodb.md)
Oracle (preview)LogMiner over the redo logs, mined from CDB$ROOTpoll: each run mines [checkpoint, SCN at open] with COMMITTED_DATA_ONLY (ADR-0037)

MySQL and PostgreSQL expose the log to the client; SQL Server does not — there a server-side Agent extracts the log into change tables that rivet polls.


Permissions & prerequisites

MySQL — the binlog grants

Rivet registers as a replica and streams the binlog. Least privilege:

CREATE USER 'rivet_cdc'@'%' IDENTIFIED BY '***';

-- read the binlog stream (register as replica, COM_BINLOG_DUMP).
-- Server-wide: REPLICATION SLAVE cannot be scoped to a database/table.
GRANT REPLICATION SLAVE  ON *.* TO 'rivet_cdc'@'%';

-- read the current binlog coordinate (SHOW MASTER STATUS) when starting
-- without a checkpoint.
GRANT REPLICATION CLIENT ON *.* TO 'rivet_cdc'@'%';

-- ONLY for `--output`: rivet resolves the table's column types with
-- `SELECT * FROM <table> LIMIT 0` (metadata only, no rows). Not needed for NDJSON.
GRANT SELECT ON `app`.`orders` TO 'rivet_cdc'@'%';

FLUSH PRIVILEGES;

Server configuration (my.cnf [mysqld], or SET GLOBAL + restart where allowed):

log_bin           = ON       # binary logging on (often already on for replication/PITR)
binlog_format     = ROW      # rivet needs row images, not statements — MIXED/STATEMENT will not work
binlog_row_image  = FULL     # full before/after image — REQUIRED for the after-image / MERGE shape;
                             # MINIMAL drops unchanged columns and breaks "overwrite all columns"
server_id         = 1        # any unique id for the source; rivet uses a DIFFERENT --server-id

Notes:

  • REPLICATION SLAVE is server-wide by design. You cannot grant binlog access for one database only — the binlog is a single server-wide stream. Scope data exposure with --table (rivet filters client-side) and the SELECT grant.

  • A stale server_id collision silently kills the stream. Give rivet a --server-id no other replica uses.

  • Connect CDC directly to MySQL — not through ProxySQL / MaxScale. The binlog stream is COM_BINLOG_DUMP, a replication protocol query proxies don’t carry; the batch path can go through a pooler, CDC cannot. Rivet probes the connection and fails fast with this exact reason if it sees a proxy, so point the source at the MySQL host (the replication endpoint), not the proxy port.

  • binlog_row_image = FULL is MySQL’s default; the risk is a source that has set it to MINIMAL to shrink the binlog — that path needs the column-mask MERGE, not the simple overwrite (see Output shape).

  • Amazon RDS / Aurora MySQL: two managed-only settings, and neither is in my.cnf. Both were diagnosed the hard way on a customer replica, a day apart.

    1. Binary logging follows automated backups. With backup retention at 0 the instance runs log_bin = 0 no matter what the parameter group says, and SHOW BINARY LOGS answers ERROR 1381 (HY000): You are not using binary logging. Set backup retention above zero (this restarts the instance), then binlog_format = ROW, binlog_row_image = FULL and binlog_row_metadata = FULL in the parameter group. A read replica also needs log_replica_updates = 1 to re-log what it applies.

    2. binlog_expire_logs_seconds does not govern retention here. RDS purges a binlog as soon as the engine itself no longer needs it — typically within minutes — so a checkpoint written by one run is unreadable by the next and the resume fails with ERROR 1236. Measured: a checkpoint taken at 13:42 was already past retention at 13:59. Set the managed knob instead, sized well above the CDC cadence:

      CALL mysql.rds_set_configuration('binlog retention hours', 72);
      CALL mysql.rds_show_configuration;   -- confirm
      

    The filenames are the tell: mysql-bin-changelog.NNNNNN is RDS’s naming, so an ERROR 1236 naming one of those is this, not PURGE BINARY LOGS.

PostgreSQL — the logical slot

Rivet’s PostgreSQL reader consumes a logical slot through the normal SQL connection (pg_logical_slot_peek_changes()), not the streaming-replication protocol. That changes what you must grant:

-- REPLICATION attribute: required to create and read a logical slot, even via
-- the SQL functions (pg_create_logical_replication_slot / _get_changes).
ALTER ROLE rivet_cdc WITH LOGIN REPLICATION PASSWORD '***';

-- ONLY for `--output`: schema resolution (SELECT ... LIMIT 0).
GRANT SELECT ON app.orders TO rivet_cdc;

Server configuration (postgresql.conf, needs a restart):

wal_level            = logical   # log enough to decode row changes (default is 'replica')
max_replication_slots = 10       # >= 1 (defaults are usually fine)
max_wal_senders       = 10       # >= 1

Notes:

  • No pg_hba.conf replication line is required. That entry is for the streaming walsender protocol; rivet’s poll model uses an ordinary connection, so the normal host app rivet_cdc ... rule suffices. (This is the main way the poll model is operationally lighter than streaming CDC tools.)
  • A logical slot pins WAL until it is consumed/advanced. An abandoned slot prevents WAL recycling and fills the disk — the number-one PostgreSQL CDC foot-gun. Drop unused slots with SELECT pg_drop_replication_slot('rivet_slot');.
  • wal2json / pgoutput are alternatives to test_decoding; test_decoding is always built in and needs no extension.

SQL Server — CDC change tables

SQL Server has no client-streamable log. A server-side Agent job extracts the log into cdc.* change tables, which rivet polls. Two distinct privilege levels:

-- ONE-TIME ENABLE (requires sysadmin or db_owner):
EXEC sys.sp_cdc_enable_db;                         -- creates the cdc schema + capture job
EXEC sys.sp_cdc_enable_table
     @source_schema = N'dbo', @source_name = N'orders',
     @role_name = N'cdc_reader',                   -- gating role for readers (or NULL = no gate)
     @capture_instance = N'dbo_orders',
     @supports_net_changes = 0;

-- RUNTIME READER (what rivet connects as — least privilege):
CREATE USER rivet_cdc FOR LOGIN rivet_cdc;
GRANT SELECT ON SCHEMA::cdc TO rivet_cdc;          -- read the change tables + functions
ALTER ROLE cdc_reader ADD MEMBER rivet_cdc;        -- if a gating role was set above

Notes:

  • SQL Server Agent must be running. The capture job (default ~5 s scan cycle) is what populates the change tables. If the Agent stops, the change tables silently freeze and the transaction log can’t truncate — disk pressure. A production reader should watch for a non-advancing sys.fn_cdc_get_max_lsn(), not read “no rows” as “no changes”.
  • Edition gate: CDC is on Enterprise / Standard / Developer — not Express or Web. On Express, use Change Tracking instead (different, lighter, but only tells you which rows changed, not the data).
  • Enabling CDC needs sysadmin/db_owner; the runtime reader needs only the SELECT grant above.
  • Retention: the cleanup job keeps ~3 days by default. If rivet is offline longer than retention, the saved LSN falls below sys.fn_cdc_get_min_lsn() and the read errors — fall back to a full re-snapshot.

Oracle — LogMiner (preview)

Oracle CDC reads the redo logs through LogMiner, which ships with every edition (Free included) and needs no GoldenGate licence. rivet never sets ENABLE_GOLDENGATE_REPLICATION (that one does need the licence).

-- ONCE, as SYSDBA in CDB$ROOT:
SHUTDOWN IMMEDIATE; STARTUP MOUNT; ALTER DATABASE ARCHIVELOG; ALTER DATABASE OPEN;
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA;              -- minimal logging
CREATE USER c##rivetcdc IDENTIFIED BY … CONTAINER = ALL;
GRANT CREATE SESSION, SET CONTAINER, LOGMINING TO c##rivetcdc CONTAINER = ALL;
GRANT EXECUTE_CATALOG_ROLE TO c##rivetcdc CONTAINER = ALL;
GRANT SELECT ON v_$database TO c##rivetcdc CONTAINER = ALL;   -- and the same for
--   v_$archived_log, v_$log, v_$logfile, v_$logmnr_contents, v_$logmnr_logs, v_$transaction

-- PER CAPTURED TABLE (its owner or a DBA), in the pluggable database:
ALTER TABLE app.orders ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;
GRANT SELECT ON app.orders TO c##rivetcdc;

Notes:

  • The URL names the pluggable database (oracle://c%23%23rivetcdc:…@host:1521/ORCLPDB1 — # percent-encoded). rivet switches the session to CDB$ROOT to mine and keeps only that PDB’s changes. A non-CDB works without the switch.
  • ALL COLUMNS logging, not PRIMARY KEY. With key logging an UPDATE’s redo carries only the changed columns, so the change cannot represent the row; rivet refuses such a table and prints the statement above.
  • Tables the preview refuses by name: a column of type LOB, LONG, XMLTYPE, JSON, INTERVAL, BOOLEAN, VECTOR, ROWID or an object type; a name over 30 bytes; and, before 23ai, an identity column (LogMiner ignores those tables entirely).
  • Retention is the DBA’s, as with the binlog: nothing pins archived logs for rivet. If the checkpoint needs a log that was deleted, the run fails with a data-loss error (see Failure modes).
  • Loading an Oracle stream with rivet load is not supported in the preview: a mode: cdc Oracle export under a load: block is refused when the config is read.

Reading from a replica (no primary access)

A common real-world constraint: you’re handed a database but only a read replica, never the master. Rivet reads the log of whatever host you point source.url at — it never needs the primary specifically. Whether a replica can serve that log is an engine + replica-config question, not a rivet limitation:

enginefrom a replica?what the replica needsverified
MySQL✅ yeslog_bin = ON and log_replica_updates = ON (log_slave_updates pre-8.0.26) so the replica re-logs replicated changes into its own binlog — this is off by default: a replica applies changes but does not re-log them without it. Plus the REPLICATION SLAVE / REPLICATION CLIENT grant and a server_id distinct from both the primary and the replica. rivet refuses a replica with log_replica_updates = OFF at start, because its binlog holds none of the replicated changes.live test + release gate: capture from a re-logging replica, refusal on one that does not
PostgreSQL✅ 16+, continuous mode onlyLogical decoding on a standby is a PostgreSQL 16 feature. Run with cdc.until_current: false: the default bounded run refuses on a standby, because the position it bounds by (pg_current_wal_lsn()) does not exist during recovery. The first run creates the slot on the standby and waits until the primary logs a running-transactions snapshot (routine on a busy primary; SELECT pg_log_standby_snapshot() on the primary forces one). Set hot_standby_feedback = on on the standby so the primary keeps the rows the slot still needs. Below 16 a standby cannot host a logical slot — point rivet at the primary.live test + release gate: continuous capture from a 16 standby, refusal of the bounded mode
SQL Server✅ yes (readable secondary)CDC is enabled and captured on the primary (the capture job runs there); the cdc.* change tables replicate to an Always On secondary with SECONDARY_ROLE (ALLOW_CONNECTIONS = ALL), and rivet reads them there with plain SELECTs.live test + release gate: a read-scale availability group (CLUSTER_TYPE = NONE), capture read from the secondary
MongoDB✅ yes (secondary)Point source.url at a secondary with readPreference=secondary (and directConnection=true for one member); the change stream reads that member’s oplog.live test + release gate: a two-member replica set, capture from the secondary

MySQL caveat — the checkpoint is replica-local. Rivet resumes by binlog {file, pos} (not GTID), and a replica’s binlog coordinates are its own, not the primary’s. A checkpoint taken against one replica does not transfer to another host, and rivet refuses one written by a different server (it records server_uuid). If you fail over (to a different replica, or to the primary), delete the checkpoint so CDC anchors on the new host first, then re-snapshot the table (mode: full).

SQL Server — the checkpoint follows a failover, and nothing else. An availability group’s replicas share one log, so a checkpoint written on the primary resumes on a secondary (verified live). The checkpoint records the database’s family_guid and recovery_fork_guid, and rivet refuses to resume against a database whose family_guid differs (another server’s database: its LSNs address a different log) or whose recovery_fork_guid changed (a RESTORE rewound the log). Recover by deleting the checkpoint so CDC re-anchors first, then re-snapshot the table with mode: full — in the other order, the changes between the snapshot and the new anchor land in neither. A checkpoint written before rivet recorded the identity resumes with a warning.

So the answer to “can I read the log from a slave?” is yes on all four engines, each verified live: MySQL (with log_replica_updates = ON), PostgreSQL 16+ in continuous mode, a SQL Server readable secondary, and a MongoDB secondary. Point source.url at the replica; everything else (grants, mode: cdc, output) is identical to running against a primary.


Output shape

--output writes one row per change in the typed after-image (upsert) shape:

__op     __pos                              __seq  id   name    amount
insert   {"file":"binlog.000046","pos":681} 0      1    alice   100
update   {"file":"binlog.000046","pos":682} 0      1    alice   150
delete   {"file":"binlog.000046","pos":683} 0      2    bob     200
  • __op — insert / update / delete.
  • __pos — the transaction’s commit position (the same value rivet checkpoints). Every change of one transaction shares it — it is not a total order over changes.
  • __seq — the change’s ordinal within its transaction (0-based, log order), so (__pos, __seq) is a total order. It is the tiebreak when one transaction touches a key more than once and, being a column, survives the load into an (unordered) warehouse table (Parquet row order does not). See CDC change ordering. (MongoDB gives every event a distinct __pos, so its __seq is always 0.)
  • the source columns, typed (resolved from the source schema), carrying the after-image for insert/update and the key (before-image) for delete.

Downstream applies it by primary key:

MERGE target t USING staged s ON t.id = s.id
WHEN MATCHED AND s.__op = 'delete' THEN DELETE
WHEN MATCHED                       THEN UPDATE SET t.* = s.*   -- overwrite all columns
WHEN NOT MATCHED AND s.__op <> 'delete' THEN INSERT (...);

With a full row image, which columns changed is irrelevant — the latest image per key already contains every prior change, so dedup-by-key + overwrite is correct. “Latest” is the highest (__pos, __seq) — the commit position, then the intra-transaction ordinal (a transaction that updates one key twice shares __pos, so __seq breaks the tie). This is why binlog_row_image = FULL matters.

Deduplicating by position, per engine

__pos granularity differs by engine, and the dedup recipe follows from it:

engine__posunique per event?
MySQL{file, pos} — the transaction’s commit positionper statement under autocommit; all events of one multi-statement transaction share it
PostgreSQL{lsn} — the transaction’s COMMIT LSNall events of one transaction share it
SQL Server{lsn} — __$start_lsnper transaction (rows within share it)

Two distinct problems:

  1. At-least-once re-delivery (a crashed run’s part re-read on resume): the re-delivered event is byte-identical — same __op, __pos, and image — so SELECT DISTINCT over the staged rows (or dedup on (pk, __pos, __op)) removes it exactly.
  2. Latest-image-per-key (the MERGE): order by (__pos, __seq) — the commit position, then the intra-transaction ordinal. __seq breaks the tie when a transaction touches one key more than once (those changes share __pos); it is a column, so — unlike Parquet row order — it survives the load. See CDC change ordering.
-- DuckDB replay: newest surviving image per key.
WITH ev AS (
  SELECT *,
         upper(lpad(split_part(__pos->>'lsn', '/', 1), 8, '0')) ||
         upper(lpad(split_part(__pos->>'lsn', '/', 2), 8, '0')) AS lsn_key   -- PostgreSQL X/Y → sortable
  FROM read_parquet('…/sessions/cdc-*.parquet')
), latest AS (
  SELECT * FROM (
    SELECT *, row_number() OVER (
      PARTITION BY id
      ORDER BY lsn_key DESC, __seq DESC) AS rn        -- (__pos, __seq) = total order
    FROM ev)
  WHERE rn = 1
)
SELECT * FROM latest WHERE __op <> 'delete';

(MySQL: order by (file, pos) parsed from __pos, then __seq; SQL Server: the fixed-width hex lsn string is already lexically ordered, then __seq.) In a warehouse MERGE, apply the same window to the staged batch first, then merge the winners by PK + __op.

Downstream loading

CDC output is the same typed Parquet the batch export writes (same build_arrow_field pipeline), so the warehouse-loading recipes apply unchanged — the engine-specific MERGE and the JSON-as-BYTES / naive-timestamp autoload recovery are in recipes/idempotent-warehouse-load.md (BigQuery) and recipes/snowflake-load.md, keyed on the PK + __op above.

Verified cross-engine on a CDC part: DuckDB reads json natively, ClickHouse as String (JSONExtract* parses it), BigQuery as BYTES (PARSE_JSON after SAFE_CONVERT_BYTES_TO_STRING); integers keep their width (INT32/INT64) and timestamps their microseconds. The JSON text round-trips losslessly in all three — the type that auto-detects differs, the data does not.

Part naming. Parts are run-stamped — cdc-<run_id>-000000.parquet — so a scheduler re-running into the same prefix appends each cycle’s parts alongside the previous cycle’s (nothing is overwritten). manifest.json / _SUCCESS describe the latest run only; a glob reader over the prefix sees the union of all cycles, which is the intended at-least-once stream — dedupe by PK + __op + __pos downstream, and archive parts you have already loaded if you want the prefix to stay small.

Without --output, rivet emits the same information as NDJSON (one JSON object per change) to stdout.

Why CDC is gentle on the source

batch (SELECT)CDC (log)
touchesthe tablethe log only
locks / read snapshotyesno
buffer-pool evictionyes (scans cold pages)no
cost scales withtable size (re-scan)change rate (deltas)
when it costsactively, every runlatently (log retention / disk)

The log is written anyway (WAL for durability, binlog for replication/PITR), so on MySQL/PostgreSQL CDC mostly reads what already exists — near-zero incremental OLTP cost. The one real CDC cost is disk via log retention if the consumer lags (PG slots pin WAL; MySQL keeps binlog until read). SQL Server is the exception: its Agent writes changes into change tables (extra write volume + storage), so CDC there trades read-contention for an ongoing write/storage overhead.

Failure modes & recovery

Every CDC run is bounded and resumable, and the durable sequence is flush → checkpoint → ack: the resume position only advances after the part is durably written. So on any failure — a dropped connection, a query error, a full source disk — the run fails loudly (non-zero exit, with the per-engine setup hint), the checkpoint/slot is not advanced, and the next run re-reads from the last good position. Rivet never silently loses a change; the trade-off is at-least-once, so a failed run’s already-uploaded parts can reappear — dedupe downstream by primary key + __op (the output is the upsert / after-image shape).

A failed run leaves its durable parts in the destination but no manifest.json / _SUCCESS — that pair marks a clean end, so a missing _SUCCESS is how you (and rivet validate) tell a partial run from a complete one.

PostgreSQL — the slot fills / the source disk fills

A logical slot pins WAL until rivet advances it (confirmed_flush_lsn). The behaviour depends on whether rivet is running:

  • Running + advancing — each successful run reads the changes, writes them durably, then advances the slot, so PostgreSQL releases the WAL up to that point. The slot only ever holds the WAL since the last advance — it does not grow unbounded while rivet keeps the slot moving.
  • Stopped (abandoned slot) — rivet does nothing (it isn’t running); the slot keeps pinning WAL and the source disk fills. This is the number-one PostgreSQL CDC foot-gun, and it is operator responsibility: SELECT pg_drop_replication_slot('rivet_slot'); when you stop capturing for good.
  • Source disk already full — run rivet (it reads WAL to advance the slot, which releases WAL and relieves the pressure) or drop the slot. If PostgreSQL is too degraded to answer, rivet’s query fails → the run fails → re-read next run.

Memory is O(largest transaction). The adapters buffer a whole transaction until its COMMIT (parts never split a transaction — the resume invariant; SQL Server buffers per poll batch). Measured on MySQL: ~1.4 KB of RSS per buffered row (~14× a 100-byte payload): a 100k-row transaction drains at ~170 MB RSS, 300k at ~440 MB — linear. A transaction past the hard caps (5M buffered rows or 2 GiB estimated bytes by default; RIVET_CDC_MAX_TX_ROWS / RIVET_CDC_MAX_TX_BYTES override them) fails the run LOUDLY before it can OOM. Opt-in RIVET_CDC_SPILL_DIR spills the adapter’s copy past the cap to disk (PostgreSQL, MySQL, SQL Server; Oracle always refuses at the cap), but the sink still holds the whole transaction, so it saves only ~11% of peak RSS — see CDC failure modes. Run bulk backfills in batched transactions, or through mode: full/initial: snapshot (the batch path streams).

DDL inside a capture window: safe where the engine names its columns, a LOUD ERROR where it does not. PostgreSQL (wire text) and SQL Server (change tables) always name every image column, so rivet maps values by NAME: a DROP COLUMN or RENAME landing between runs captures correctly, and an equal-arity DROP a + ADD c leaves c NULL for the older images rather than filling it with a neighbour’s value (unless the dropped column sat at c’s position, which looks exactly like a rename and is read as one). A column ADDED while a run is open is not in that run’s schema, so its values for that run’s window are dropped — re-snapshot the table after an ADD COLUMN if those values matter. MySQL’s binlog carries names only when the server runs with binlog_row_metadata=FULL (8.0.1+ — strongly recommended; the compose test stack sets it):

# my.cnf — makes mid-stream DDL safe for rivet CDC
binlog_row_metadata = FULL

Under the default MINIMAL the binlog is nameless and positional — expect runs to FAIL with an explicit error (“an event … carries N column(s) but the resolved schema has M”) whenever a DDL lands inside a capture window. That is deliberate: mapping by position would put values into the wrong columns silently, and a loud stop is the only safe behavior. Recover by re-snapshotting the table (or resetting the checkpoint past the DDL), and set binlog_row_metadata=FULL to retire this error class. DDL between runs is always fine — each run resolves the schema fresh. A mid-window RENAME is safe in both modes (same arity ⇒ positional fallback keeps the value). Same-arity TYPE changes remain undetectable without schema history (roadmap) — run type migrations and their backfills through a re-snapshot.

The value checksum runs on CDC too. The same always-on two-ended check the batch export performs — an independent fold of the decoded cells vs a fold of the built Arrow column — runs per column before every CDC part is written; a mismatch fails the run naming the column, never writes the corrupted part. Failure behaviour is also parity: a value unrepresentable in the declared column (PostgreSQL 'NaN'::numeric in a Parquet decimal) fails loudly on both paths, never a silent NULL.

For the full operational failure playbook — every symptom, what rivet does, how to recover, how to prevent — see cdc-failure-modes.md.

A vanished slot is a loud error, not a silent restart. When a resume checkpoint exists but the slot is gone (dropped by an operator, or invalidated and removed), rivet refuses to re-create it — a fresh slot would anchor at the current position and silently skip everything since the drop. The run fails with the re-snapshot hint; delete the checkpoint file only when you explicitly accept a fresh anchor.

Bound the blast radius: set max_slot_wal_keep_size (PG 13+). PostgreSQL then invalidates the slot rather than fill the disk; rivet’s next run fails with a slot-invalidated error and you re-snapshot. Monitor pg_replication_slots (active, and restart_lsn vs the current LSN = how much WAL the slot is holding).

rivet doctor automates this monitoring. For a config with mode: cdc exports, doctor probes the engine: PostgreSQL — the export’s slot (retained WAL, fails above 1 GiB) and any other inactive slot pinning WAL (the abandoned-slot foot-gun); MySQL — log_bin/binlog_format=ROW/ binlog_row_image=FULL, and whether the checkpoint’s binlog file is still retained (a purged file is reported before the run fails with ERROR 1236); SQL Server — CDC enabled, the capture instance exists, the checkpoint is within retention, and the Agent service is running.

MySQL — the binlog was purged

If rivet is offline long enough that the saved binlog position is purged (binlog_expire_logs_seconds / PURGE BINARY LOGS), the resume read fails with MySQL ERROR 1236 (the requested binlog file is gone). The position is unrecoverable — delete the checkpoint so CDC re-anchors first, then re-snapshot (mode: full). Size binlog retention comfortably above your CDC cadence.

SQL Server — the checkpoint fell below retention

If the saved LSN falls below sys.fn_cdc_get_min_lsn() (the cleanup job — ~3 days by default — removed the changes after it), rivet fails loudly — “the resume position is older than the change-table retention … re-snapshot” — rather than resume from the new min and silently skip the gap. Delete the checkpoint so CDC re-anchors first, then re-snapshot. Also watch for a non-advancing sys.fn_cdc_get_max_lsn(): that means the Agent capture job stopped, so the change tables are frozen — read “no rows” as “the job is down”, not “no changes”.

Oracle — the archived logs were deleted

If the checkpoint needs redo older than the oldest archived log still listed (RMAN DELETE INPUT, a retention policy), or a log sequence is missing in between, the run fails with “… LOST to this stream” instead of mining from whatever remains. Delete the checkpoint so the next run anchors first, then re-snapshot the table.

A checkpoint written against another database (a different DBID, a RESETLOGS since, or another pluggable database) is refused the same way: an SCN means nothing outside the database that issued it.

Recovery, in one line

Re-run to resume from the last checkpoint (the common case). If the run reports the position is unrecoverable (PostgreSQL slot invalidated, MySQL binlog purged, SQL Server retention exceeded), restart CDC from a new checkpoint first, then re-snapshot the table with mode: full — the only safe recovery once the source log no longer covers the gap. The order matters: the new anchor must exist before the snapshot reads, so the stream overlaps the snapshot (duplicates, which the load deduplicates) instead of leaving the changes in between in neither.

Limitations (current)

Typed output (real Timestamp/Date32/Decimal128), commit-boundary checkpointing, cloud destinations + manifest.json/_SUCCESS, and the config-driven rivet run path with a recorded run are all in place for all three engines. What remains:

  • Continuous capture is bounded-poll-and-exit by default on every engine (they drain their backlog and stop). The supported continuous model is a scheduler running the default bounded rivet cdc (or rivet run with cdc.until_current: true, now the default) on an interval, each run resuming from the checkpoint. For an unbounded run, pass --stream (config: cdc.until_current: false; the config-driven rivet run path logs an engine-specific warning), because only MySQL (the binlog dump blocks) and MongoDB (the change stream blocks awaiting events, ending only if the stream is invalidated or closed) genuinely stay up as daemons; PostgreSQL and SQL Server still exit on catch-up (one unbounded pass — run it under a supervisor that restarts it). The bounded run remains the intended model.
  • Schema drift: the sink schema is frozen at the first flush — a column added mid-run is not picked up until the next run re-resolves the table, and its values captured in the meantime are dropped (the events are still acked) — re-snapshot the table to recover them.
  • No lag metric: the run records rows / files / bytes / duration / status, but not replication lag (“how far behind the source is”) — the next observability step.
  • Pre-image completeness depends on the source config: full UPDATE/DELETE before-images need binlog_row_image=FULL (MySQL) / REPLICA IDENTITY FULL (PostgreSQL); otherwise only key columns are carried.
  • Type parity with the batch export is total: every Rivet-mapped type — including PostgreSQL arrays (real List columns, inner NULLs preserved) and NUMERIC/DECIMAL above precision 38 (Decimal256) — is byte-identical to the batch export, enforced per engine by the live *_full_type_matrix_matches_batch tests (ArrayData equality).