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

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.

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 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 and POST it to the Connect REST API endpoint:

{
  "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, 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 file demonstrates a Flink job that consumes CDC events, applies windowing to collapse duplicates, and writes to Delta Lake:

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 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.
  • Schema evolution is managed through Confluent Schema Registry, as recommended in 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 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 pattern recommends using Flink's tumbling windows or similar techniques to collapse duplicate events that occur due to retries or reprocessing.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →