Debezium tails the log

Debezium is a Kafka Connect connector that subscribes to the WAL/binlog over the streaming replication protocol — not SQL — and decodes each record into a change event {op, before, after, source} emitted as one Kafka message per row change.

Previously

The DB already has a durable, ordered log of every change — so the next move is a process that tails it and turns each record into a downstream event.

Scene 04

Debezium tails the log

  1. Watch
  2. Try it
  3. Predict
  4. Capture
DB · Postgrestablesid=42 │ status=pending │ amount=$120id=73 │ status=shipped │ amount=$48id=08 │ status=pending │ amount=$305WALstreaming replication protocolDebezium connector(stopped)(waiting for first WAL record)Kafka topic(no events yet)
What to watch for

We need a process that tails the WAL and emits one event per row change. Debezium is the canonical open-source implementation; it runs as a Kafka Connect connector and emits one change event per row change. Watch: an UPDATE lands in the WAL, the connector decodes it, the topic gets one message.

Continue unlocks when the animation finishes.
Implementation

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

KafkaConnect.subscribeAsReplica
boot path: register as a streaming-replication subscriber
def start(connectorConfig):
conn = postgres.connect(replication='database')
conn.create_replication_slot(
name = connectorConfig.slot_name,
plugin = 'pgoutput', # logical decoding
)
stream = conn.start_replication(
slot = connectorConfig.slot_name,
start_lsn = offsets.last_confirmed_lsn(),
)
for rec in stream:
Connector.onWalRecord(rec)
Connector.onWalRecord
main loop: decode one WAL record, emit one Kafka message
def onWalRecord(rec):
if connector.mode == 'select':
# SELECT path doesn't see deletes
rows = db.query('SELECT * FROM ' + rec.table)
return emitRowDiffs(rows)
# streaming-replication path: decode the WAL tuple
before, after = decodeTuples(rec)
op = {INSERT:'c', UPDATE:'u', DELETE:'d'}[rec.kind]
envelope = buildEnvelope(rec, op, before, after)
kafka.produce(
topic = topicFor(rec.table),
key = pk(after or before),
value = envelope,
)
buildEnvelope
construct the {op, before, after, source} change event
def buildEnvelope(rec, op, before, after):
return {
'op': op, # 'c' | 'u' | 'd'
'before': before, # null on insert
'after': after, # null on delete
'source': {
'table': rec.table,
'lsn': rec.lsn,
'ts_ms': rec.commit_ts_ms,
},
}

Where this sits in Build a CDC pipeline (Debezium + outbox)

Scene 04 of 12. Debezium is a connector that registers as a replica, decodes each WAL/binlog record into a structured change event, and emits one Kafka message per row change.

Up next. Debezium tails the log and emits one change event per row change — but if it crashes, the database needs to know where it left off so the WAL can be safely recycled.

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