Snapshot, then stream
Debezium runs a consistent snapshot of the table emitted as op=r events, records the LSN at snapshot time, then streams from that LSN — so history and live updates are stitched at a single seam, and rows touched mid-snapshot show up twice.
The slot tracks where the connector left off in the WAL — but the WAL only goes back so far, so a connector attaching to a database with a billion existing rows cannot replay history from the log alone. So how does Debezium hand a fresh consumer the entire table without the log to lean on?
Scene 06
Snapshot, then stream
- Watch
- Try it
- Predict
- Capture
100M-row accounts table. Watch the snapshot bar drain emitting op=r events tagged with the snapshot LSN. Concurrent writes accumulate in the thin lane above the WAL strip. Once the bar empties, the seam is crossed and live op=u events from the WAL begin to flow.
Highlighted lines are the ones running in the diagram right now.
def initialSnapshot():if snapshot.mode == 'never':return streamFromLsn(currentLsn()) # skip historytx = db.begin('REPEATABLE READ')snapshotLsn = tx.currentLsn()if snapshot.mode != 'no_data':for table in captured.tables:for row in tx.select_all(table):emit(op='r', after=row, lsn=snapshotLsn)tx.commit()streamFromLsn(snapshotLsn) # seam
def streamFromLsn(lsn):cursor = replicationSlot.startReplication(startLsn = lsn,)for change in cursor:# change.lsn >= snapshotLsn — concurrent# writes during snapshot replay here tooemit(op = change.op, # c | u | dbefore = change.before,after = change.after,lsn = change.lsn,)
# A row updated WHILE the snapshot was running surfaces twice:# 1. op=r tagged at snapshotLsn (from the bulk read)# 2. op=u tagged at change.lsn > snapshotLsn (from the WAL)# The connector cannot suppress (2) without risking a dropped# write, so dedupe is the consumer's job.def onEvent(evt):key = (evt.table, evt.pk)if evt.lsn <= seen.get(key, -1):return # already applied — dropapply(evt)seen[key] = evt.lsn
Where this sits in Build a CDC pipeline (Debezium + outbox)
Scene 06 of 12. Debezium runs a consistent snapshot (op=r events), records the LSN at snapshot time, then switches to streaming from that LSN — so history and live updates stitch at one seam.
Up next. Now the pipeline can faithfully replay history and stream live changes — but it is emitting raw row deltas with whatever column names the DBA picked, which means the next ALTER TABLE leaks straight into every consumer.
All 12 scenes in Build a CDC pipeline (Debezium + outbox) · Every curriculum