# How to Implement Change Data Capture (CDC) Patterns with Kafka Connect

> Learn to implement Change Data Capture CDC patterns with Kafka Connect using Debezium and Apache Flink for real-time data streaming and analytics. Stream database changes effortlessly.

- Repository: [DataExpert.io/data-engineer-handbook](https://github.com/DataExpert-io/data-engineer-handbook)
- Tags: how-to-guide
- Published: 2026-08-07

---

**Implement Change Data Capture (CDC) patterns with Kafka Connect by deploying a source connector like Debezium to stream database change logs into Kafka topics, then consume those events with stream processors such as Apache Flink for real-time analytics.**

Change Data Capture (CDC) enables streaming analytics by capturing row-level changes in source databases and propagating them to downstream systems in near real-time. According to the DataExpert-io/data-engineer-handbook repository, Kafka Connect serves as the backbone for production CDC pipelines, providing a fault-tolerant framework for reading change logs and writing them to Kafka topics without custom code.

## Understanding the CDC Architecture with Kafka Connect

The handbook outlines a five-layer architecture for CDC pipelines in [`intermediate-bootcamp/introduction.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/introduction.md).

### Source Database Layer

Enable the database's native change-log mechanism—such as MySQL binlog, PostgreSQL WAL, or CDC-enabled streams—to publish row-level modifications.

### Kafka Connect Source Connector

Deploy a source connector (e.g., Debezium) that reads the change log and emits events to Kafka. The connector handles initial snapshots and ongoing streaming.

### Kafka Topics and Schema Registry

Raw change events land in Kafka topics preserving source order. The [`design-patterns/data-developer-platform.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/data-developer-platform.md) file recommends integrating Confluent Schema Registry to manage Avro or JSON schemas, enabling safe schema evolution without breaking consumers.

### Stream Processors and Consumers

Downstream systems like Apache Flink, Spark Structured Streaming, or simple Kafka consumers read from topics to apply transformations, joins, or aggregations before loading into data lakes or warehouses.

## Deploying a Debezium MySQL CDC Connector

To implement CDC with Kafka Connect, configure a source connector using the Debezium MySQL connector class. Save the following configuration as [`mysql-cdc-connector.json`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/mysql-cdc-connector.json) and POST it to the Connect REST API endpoint:

```json
{
  "name": "mysql-cdc-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "mysql-host",
    "database.port": "3306",
    "database.user": "debezium",
    "database.password": "debezium-pw",
    "database.server.id": "184054",
    "database.server.name": "mydb",
    "database.whitelist": "inventory",
    "database.history.kafka.bootstrap.servers": "kafka:9092",
    "database.history.kafka.topic": "schema-changes.inventory",
    "include.schema.changes": "true",
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.drop.tombstones": "true",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false"
  }
}

```

Key configuration details include setting `transforms` to `unwrap` with the `ExtractNewRecordState` type, which flattens the Debezium envelope to extract only the new record state. The `database.history.kafka.topic` parameter ensures schema changes are tracked in a dedicated Kafka topic.

## Ensuring Exactly-Once Semantics and Deduplication

According to [`design-patterns/microbatch-deduplication.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/microbatch-deduplication.md), CDC pipelines must handle duplicate events that arise from connector retries or consumer rebalances. Configure the connector with `consumer.auto.offset.reset=earliest` and implement micro-batch deduplication in your stream processor.

The [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py) file demonstrates a Flink job that consumes CDC events, applies windowing to collapse duplicates, and writes to Delta Lake:

```python
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment, DataTypes

env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)

# Define the source table that reads from Kafka

t_env.execute_sql("""
CREATE TABLE mysql_cdc (
  id BIGINT,
  name STRING,
  updated_at TIMESTAMP(3),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'kafka',
  'topic' = 'mydb.inventory.customers',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'json',
  'json.fail-on-missing-field' = 'false'
)
""")

# Apply a 5‑second tumbling window to deduplicate

t_env.execute_sql("""
INSERT INTO delta_table
SELECT id, LAST_VALUE(name) AS name, MAX(updated_at) AS updated_at
FROM mysql_cdc
WINDOW TUMBLE (5 SECONDS)
GROUP BY id
""")

```

This implementation uses a 5-second tumbling window to deduplicate records, ensuring exactly-once semantics when writing to the destination table.

## Managing Schema Evolution

The [`design-patterns/data-developer-platform.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/data-developer-platform.md) pattern emphasizes using a Schema Registry to handle evolving database schemas. When source tables add columns or change data types, the Debezium connector captures these schema changes in the history topic. Consumers configured with schema registry clients can automatically adapt to new schemas without job restarts.

## Summary

- **Kafka Connect** provides the framework for CDC by reading database change logs and publishing them to Kafka topics using connectors like Debezium.
- **Configuration** requires setting `transforms.unwrap.type` to `ExtractNewRecordState` and enabling schema history tracking via `database.history.kafka.topic`.
- **Exactly-once processing** combines connector-level settings with stream processing windowing techniques documented in [`design-patterns/microbatch-deduplication.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/microbatch-deduplication.md).
- **Schema evolution** is managed through Confluent Schema Registry, as recommended in [`design-patterns/data-developer-platform.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/data-developer-platform.md), ensuring downstream consumers handle structural changes gracefully.
- **End-to-end implementation** examples are available in [`intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/intermediate-bootcamp/materials/4-apache-flink-training/src/job/start_job.py) for consuming CDC streams with Apache Flink.

## Frequently Asked Questions

### What is the primary benefit of using Kafka Connect for CDC?

Kafka Connect provides a standardized, scalable framework for capturing database changes without writing custom polling scripts. It handles fault tolerance, offset management, and schema tracking automatically, enabling near real-time data propagation to Kafka topics.

### How does Debezium handle schema changes in the source database?

Debezium captures DDL events in the `database.history.kafka.topic` and emits them alongside DML changes. When `include.schema.changes` is set to `true`, downstream consumers receive schema update events that can be processed to adjust target table structures accordingly.

### What is the purpose of the unwrap transformation in CDC connectors?

The `unwrap` transformation using `io.debezium.transforms.ExtractNewRecordState` removes the Debezium envelope metadata, exposing only the new row state. This simplifies downstream processing by eliminating the need to parse complex nested structures containing both before and after images.

### How can I prevent duplicate records when processing CDC events?

Configure the Kafka Connect source with `consumer.auto.offset.reset=earliest` and implement windowed deduplication in your stream processor. The handbook's [`design-patterns/microbatch-deduplication.md`](https://github.com/DataExpert-io/data-engineer-handbook/blob/main/design-patterns/microbatch-deduplication.md) pattern recommends using Flink's tumbling windows or similar techniques to collapse duplicate events that occur due to retries or reprocessing.