1. Quick Start: Dựng CDC pipeline MySQL với Debezium và Kafka trong 5 phút
Trước đây, để đồng bộ dữ liệu sang Elasticsearch hay các microservice khác, team mình quen tay setup cronjob quét theo updated_at cứ mỗi 5 phút. Giải pháp này chạy ổn cho đến khi bảng users chạm mốc 12 triệu dòng. Lúc đó, DB liên tục nghẽn I/O, CPU vọt lên 90% mỗi lần cronjob quét qua. CDC (Change Data Capture) sinh ra để giải quyết triệt để vấn đề này. Thay vì query trực tiếp vào bảng, ta đọc thẳng stream từ binary log (binlog) của MySQL.
Dưới đây là file docker-compose.yml hoàn chỉnh gồm MySQL 8.0, Zookeeper, Apache Kafka và 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
Khởi động cụm service:
docker compose up -d
Chờ khoảng 30 giây cho Kafka Connect khởi động xong, bạn đăng ký MySQL Connector qua REST API tại 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"
}
}'
Tạo thử bảng và chèn bản ghi kiểm tra:
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]');
Mở terminal lắng nghe Kafka topic để thấy event CDC bắn về theo thời gian thực (độ trễ thường dưới 15ms):
docker compose exec kafka /kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka:9092 \
--topic dbserver1.inventory.users \
--from-beginning
2. Cơ chế hoạt động đằng sau
2.1. Debezium đọc Binlog như thế nào?
Bản chất Debezium đóng vai trò như một MySQL Replica. Nó kết nối tới server, xin đọc stream binlog liên tục và phân tích payload binary thành JSON có cấu trúc. Mọi thao tác ghi dữ liệu đều biến thành event gửi thẳng vào Kafka topic tương ứng.
2.2. 3 cờ cấu hình binlog bắt buộc phải nhớ
- binlog-format=ROW: Bắt buộc Debezium phân tích dữ liệu ở cấp độ dòng thay vì câu lệnh. Nếu để
STATEMENT, log chỉ lưu cú pháp SQL thô và Debezium không thể biết chính xác giá trị các cột thay đổi ra sao. - binlog-row-image=FULL: MySQL sẽ ghi lại đầy đủ toàn bộ trường dữ liệu trước (before) và sau (after) khi sửa. Nhờ vậy, consumer phía sau nhận trọn vẹn context mà không phải query ngược lại DB gốc.
- gtid-mode=ON: Bật mã định danh giao dịch toàn cục. Khi hệ thống xảy ra failover giữa Master và Replica, Debezium vẫn giữ đúng offset binlog mà không bị mất hay duplicate event.
2.3. Cấu trúc Message JSON trong Kafka
Mỗi message trên Kafka gồm 2 phần: schema và payload. Trừ khi dùng Avro/Protobuf để nén schema, bạn sẽ làm việc trực tiếp với các trường trong payload:
before: Dữ liệu của dòng trước khi câu lệnh thực thi (luôn trả vềnullvới lệnh INSERT).after: Dữ liệu mới của dòng sau khi thay đổi (trả vềnullvới lệnh DELETE).op: Mã thao tác (c: Create/Insert,u: Update,d: Delete,r: Read snapshot ban đầu).ts_ms: Timestamp tính bằng epoch millisecond khi event được commit tại MySQL.
3. Thiết lập nâng cao cho môi trường Production
3.1. Tạo dedicated user cho Debezium
Không bao giờ cấp quyền root cho connector trên production. Bạn nên tạo riêng một user với nhóm quyền tối thiểu cần thiết để đọc replica:
CREATE USER 'debezium_user'@'%' IDENTIFIED BY 'MatKhauPhucTap123!@#';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium_user'@'%';
FLUSH PRIVILEGES;
3.2. Lọc bảng và cột nhạy cảm
Để tránh lọt dữ liệu nhạy cảm như password hash hay số thẻ ngân hàng vào Kafka topic, hãy dùng whitelist và blacklist ở mức connector:
{
"table.include.list": "inventory.users,inventory.orders",
"column.exclude.list": "inventory.users.password_hash,inventory.users.credit_card"
}
3.3. Tối ưu Snapshot khi khởi động trên bảng dung lượng lớn
Theo mặc định, snapshot.mode = initial sẽ thực hiện SELECT toàn bộ bảng lúc bật connector. Với bảng trên 20 triệu dòng, việc này dễ gây lock bảng hoặc tràn RAM. Bạn nên cân nhắc chuyển sang Incremental Snapshot hoặc điều chỉnh kích thước batch đọc:
{
"snapshot.mode": "when_needed",
"snapshot.locking.mode": "none",
"snapshot.fetch.size": 10240
}
4. Kinh nghiệm thực chiến sau 6 tháng vận hành
4.1. Bài toán Retention Binlog: Tránh mất offset khi connector downtime
Một bài học đau thương team mình từng dính: Kafka Connect bị treo suốt 2 ngày cuối tuần do hết disk trên worker node. Đến sáng thứ 2 khôi phục lại thì MySQL đã purge hết binlog cũ. Debezium mất dấu offset và buộc phải trigger lại full snapshot từ đầu, gây nghẽn toàn bộ luồng consumer phía sau. Để an toàn, hãy đặt retention cho binlog tối thiểu từ 3 đến 7 ngày:
-- MySQL 8.0: Giữ binlog trong 7 ngày (604800 giây)
SET GLOBAL binlog_expire_logs_seconds = 604800;
4.2. Bẫy Tombstone Records khi xử lý lệnh DELETE
Mỗi khi có thao tác xoá dòng trong MySQL, Debezium gửi liên tiếp 2 message vào Kafka: một message chứa op: "d" với dữ liệu cũ, kèm theo ngay sau đó là một message rỗng mang payload null (gọi là Tombstone Record để phục vụ Kafka Log Compaction). Nếu consumer phía backend không handle trường hợp payload null, code sẽ dính lỗi NullPointerException lập tức.
4.3. Giám sát Lag thời gian thực qua Prometheus & Grafana
Đừng để consumer bị trễ dữ liệu rồi mới phát hiện. Hãy xuất metric JMX từ Debezium sang Prometheus và chú ý metric MilliSecondsBehindSource. Nếu chỉ số này vượt quá 5000ms trong giờ cao điểm, đó là tín hiệu cảnh báo bạn cần tăng CPU/RAM cho Kafka Connect worker hoặc kiểm tra lại băng thông mạng giữa DB và Kafka.

