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.

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
Snapshot → Stream · seam at snapshot LSNmode = initial · 0 of 100,000,000 rows snapshottedsnapshot at t0 · op=r0% drainedsnapshot LSNlsn=1000concurrent writers (during snapshot)WAL · op=c/u/d (post-seam)head=lsn:1000Kafka topic(waiting)Snapshot draining — op=r events tagged at snapshotLsn.
What to watch for

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.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Connector.initialSnapshot
consistent bulk read, then handoff at snapshotLsn
def initialSnapshot():
if snapshot.mode == 'never':
return streamFromLsn(currentLsn()) # skip history
tx = 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
Connector.streamFromLsn
tail the WAL forward from the recorded seam
def streamFromLsn(lsn):
cursor = replicationSlot.startReplication(
startLsn = lsn,
)
for change in cursor:
# change.lsn >= snapshotLsn — concurrent
# writes during snapshot replay here too
emit(
op = change.op, # c | u | d
before = change.before,
after = change.after,
lsn = change.lsn,
)
Consumer.expectDuplicatesAtSeam
the dedupe contract pushed onto downstream consumers
# 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 — drop
apply(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

Built with Arqly
Every scene in Build a CDC pipeline (Debezium + outbox) builds on the one before it.All 12 Build a CDC pipeline (Debezium + outbox) scenes