How to Configure Change Data Capture (CDC) in MySQL with Debezium and Apache Kafka

MySQL tutorial - IT technology blog
MySQL tutorial - IT technology blog

1. Quick Start: Build a MySQL CDC Pipeline with Debezium and Kafka in 5 Minutes

Previously, to synchronize data to Elasticsearch or other microservices, our team used to set up cron jobs scanning by updated_at every 5 minutes. This approach worked fine until the users table reached 12 million rows. At that point, the database suffered continuous I/O bottlenecks and CPU spikes to 90% each time the cron job executed. CDC (Change Data Capture) was born to solve this exact problem. Instead of querying tables directly, we read the stream straight from MySQL’s binary log (binlog).

Below is the complete docker-compose.yml file containing MySQL 8.0, Zookeeper, Apache Kafka, and Debezium Kafka Connect:

version: '3.8'
services:
  zookeeper:
    image: quay.io/debezium/zookeeper:2.5
    ports:
      - "2181:2181"

  kafka:
    image: quay.io/debezium/kafka:2.5
    ports:
      - "9092:9092"
    links:
      - zookeeper
    environment:
      - ZOOKEEPER_CONNECT=zookeeper:2181
      - KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092

  mysql:
    image: mysql:8.0
    ports:
      - "3306:3306"
    environment:
      - MYSQL_ROOT_PASSWORD=debezium
      - MYSQL_DATABASE=inventory
      - MYSQL_USER=mysqluser
      - MYSQL_PASSWORD=mysqlpw
    command: >
      --server-id=223344
      --log-bin=mysql-bin
      --binlog-format=ROW
      --binlog-row-image=FULL
      --gtid-mode=ON
      --enforce-gtid-consistency=ON

  debezium:
    image: quay.io/debezium/connect:2.5
    ports:
      - "8083:8083"
    links:
      - kafka
      - mysql
    environment:
      - BOOTSTRAP_SERVERS=kafka:9092
      - GROUP_ID=1
      - CONFIG_STORAGE_TOPIC=my_connect_configs
      - OFFSET_STORAGE_TOPIC=my_connect_offsets
      - STATUS_STORAGE_TOPIC=my_connect_statuses

Start the service cluster:

docker compose up -d

Wait about 30 seconds for Kafka Connect to finish starting up, then register the MySQL Connector via the REST API on port 8083:

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" \
  http://localhost:8083/connectors/ -d '{
  "name": "inventory-connector",
  "config": {
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "tasks.max": "1",
    "database.hostname": "mysql",
    "database.port": "3306",
    "database.user": "root",
    "database.password": "debezium",
    "database.server.id": "184054",
    "topic.prefix": "dbserver1",
    "database.include.list": "inventory",
    "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
    "schema.history.internal.kafka.topic": "schema-changes.inventory"
  }
}'

Create a test table and insert a sample record:

CREATE TABLE inventory.users (
  id INT AUTO_INCREMENT PRIMARY KEY,
  name VARCHAR(255),
  email VARCHAR(255),
  created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);

INSERT INTO inventory.users (name, email) VALUES ('Nguyen Van A', '[email protected]');

Open a terminal to consume the Kafka topic and watch CDC events stream in real time (latency is typically under 15ms):

docker compose exec kafka /kafka/bin/kafka-console-consumer.sh \
  --bootstrap-server kafka:9092 \
  --topic dbserver1.inventory.users \
  --from-beginning

2. How It Works Under the Hood

2.1. How Does Debezium Read the Binlog?

In essence, Debezium acts as a MySQL Replica. It connects to the server, requests a continuous binlog stream, and parses the binary payload into structured JSON. Every write operation is transformed into an event and sent directly to the corresponding Kafka topic.

2.2. 3 Must-Know Binlog Configuration Flags

  • binlog-format=ROW: Forces Debezium to parse data at the row level rather than the statement level. If set to STATEMENT, the log only stores raw SQL statements, preventing Debezium from determining the exact modified column values.
  • binlog-row-image=FULL: MySQL logs all column data before and after the modification. This allows downstream consumers to receive full context without having to query the source database.
  • gtid-mode=ON: Enables Global Transaction Identifiers. When a failover occurs between Master and Replica, Debezium maintains the exact binlog offset without losing or duplicating events.

2.3. Kafka JSON Message Structure

Each message in Kafka consists of two parts: schema and payload. Unless you use Avro or Protobuf for schema compression, you will work directly with the fields inside the payload:

  • before: The row data before the statement was executed (always null for INSERT operations).
  • after: The updated row data after the change (always null for DELETE operations).
  • op: Operation type (c: Create/Insert, u: Update, d: Delete, r: Read initial snapshot).
  • ts_ms: The epoch millisecond timestamp when the event was committed in MySQL.

3. Advanced Setup for Production Environments

3.1. Create a Dedicated User for Debezium

Never grant root privileges to connectors in production. You should create a dedicated user with the minimum required permissions to read replicas:

CREATE USER 'debezium_user'@'%' IDENTIFIED BY 'StrongPassword123!@#';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium_user'@'%';
FLUSH PRIVILEGES;

3.2. Filter Sensitive Tables and Columns

To prevent sensitive data such as password hashes or credit card numbers from leaking into Kafka topics, configure whitelisting and blacklisting at the connector level:

{
  "table.include.list": "inventory.users,inventory.orders",
  "column.exclude.list": "inventory.users.password_hash,inventory.users.credit_card"
}

3.3. Optimize Startup Snapshot for Large Tables

By default, snapshot.mode = initial executes a full table SELECT when starting the connector. For tables with over 20 million rows, this can easily lead to table locks or memory exhaustion. Consider switching to incremental snapshotting or tuning the fetch batch size:

{
  "snapshot.mode": "when_needed",
  "snapshot.locking.mode": "none",
  "snapshot.fetch.size": 10240
}

4. Battle-Tested Lessons from 6 Months in Production

4.1. The Binlog Retention Challenge: Preventing Offset Loss During Connector Downtime

A painful lesson our team learned the hard way: Kafka Connect froze over an entire weekend due to disk space running out on a worker node. By Monday morning when we recovered it, MySQL had already purged the older binlogs. Debezium lost its offset position and was forced to trigger a full snapshot from scratch, causing massive bottlenecks across all downstream consumers. To be safe, set your binlog retention to at least 3 to 7 days:

-- MySQL 8.0: Keep binlog files for 7 days (604800 seconds)
SET GLOBAL binlog_expire_logs_seconds = 604800;

4.2. The Tombstone Record Trap When Handling DELETE Operations

Whenever a row deletion occurs in MySQL, Debezium sends two consecutive messages to Kafka: one containing op: "d" with the previous data, followed immediately by an empty message with a null payload (known as a Tombstone Record for Kafka Log Compaction). If your backend consumer fails to handle null payloads, it will instantly throw a NullPointerException.

4.3. Real-Time Lag Monitoring with Prometheus & Grafana

Don’t wait until consumers suffer from data delays to notice an issue. Export JMX metrics from Debezium to Prometheus and monitor MilliSecondsBehindSource closely. If this metric exceeds 5000ms during peak hours, it is a warning sign that you need to scale CPU/RAM for Kafka Connect workers or inspect network bandwidth between your DB and Kafka.

Share: