Designing a CDC Pipeline: Database to Kafka with Schema Evolution and Recovery
A CDC pipeline from a relational database into Kafka has four critical design areas: the source connector, the schema management strategy, the snapshot and offset handling, and the recovery model. The standard approach uses Debezium, which is a set of source connectors that read the database's transaction log (WAL for PostgreSQL, binlog for MySQL) and emit change events to Kafka. The source connector must have a replication slot or equivalent to hold the log position, and it must track offsets so that it can resume from the correct position after a failure. The first design decision is the snapshot mode: initial runs a full snapshot of the tables before streaming changes, never skips the snapshot and only streams new changes, and initial_only does a snapshot then stops. For a new pipeline, initial is the right default because it captures the current state before streaming changes.
Schema evolution is the second critical area. Debezium captures the database schema and emits it with each change event, but downstream consumers need a stable schema contract. The standard solution is to use a schema registry with Avro or Protobuf. The source connector writes schemas to the registry, and consumers fetch the schema by ID. For schema changes, the connector propagates DDL events. Adding a column is typically backward-compatible because consumers that do not know about the new field can ignore it. Dropping a column or changing a type is a breaking change and requires a new schema subject or a coordinated consumer upgrade. The source connector should be configured with the schema registry URL, and the schema evolution strategy should be defined before the pipeline goes live. If you are not using a schema registry, you are relying on JSON with schemas embedded, which is less robust and harder to evolve.
Recovery and offset management are the third area. Debezium stores offsets in the connect-offsets topic in distributed mode. On restart, the connector reads its offset and resumes from that position. The critical failure mode to watch is replication slot lag: if the connector stops acknowledging the slot, the database accumulates WAL and can run out of disk. Monitor slot lag and set alerts. The fourth area is the downstream contract. Consumers of the CDC topic must be prepared to handle insert, update, and delete events, and they must be prepared for schema changes. A common mistake is to treat the CDC topic as a simple stream of rows and ignore the operation type. Another mistake is to run Debezium with snapshot.mode=always on a large table, which causes a full snapshot on every restart and can flood the topic. The trade-off between snapshot modes is between initial consistency and restart cost. Version note: Debezium is an independent project, not part of Apache Kafka; its features and configuration depend on the version you deploy. The Kafka Connect framework itself is stable, but connector-specific behavior varies.
Use Debezium for CDC from relational databases; it reads the transaction log and emits change events.
Snapshot mode: initial for first deployment, never if topics already have data, always only for small tables.
Use a schema registry with Avro or Protobuf to manage schema evolution; add column is usually safe, drop or type change is breaking.
Offsets are stored in connect-offsets; on restart, the connector resumes from the last committed offset.
Monitor replication slot lag; an idle connector can cause WAL accumulation and fill the database disk.
Downstream consumers must handle insert, update, and delete events and be prepared for schema changes.
Debezium is an independent project; its behavior and config depend on the version.
0-2 years experience
2-5 years experience
5-8 years experience
8+ years experience