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.
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
- Watch
- Try it
- Predict
- Capture
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.
Highlighted lines are the ones running in the diagram right now.
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)
def onWalRecord(rec):if connector.mode == 'select':# SELECT path doesn't see deletesrows = db.query('SELECT * FROM ' + rec.table)return emitRowDiffs(rows)# streaming-replication path: decode the WAL tuplebefore, 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,)
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