Debezium 案例详解
Debezium 是一个开源的分布式平台,用于捕获数据库中的变更数据(CDC,Change Data Capture),它能够实时捕获数据库中的插入、更新和删除操作,并将这些变更流式传输到 Kafka 等消息系统中。

核心概念
graph LR
A[数据库 MySQL/PostgreSQL/MongoDB] -->|捕获变更| B[Debezium Connector]
B -->|发送到| C[Kafka Topic]
C -->|消费| D[下游应用/数据仓库/搜索引擎]
典型应用场景
-
数据库复制与迁移
- 从旧数据库实时复制到新数据库
- 构建数据湖/数据仓库的实时ETL
-
微服务数据同步
- 多个微服务共享数据但各自维护独立的数据库
- 通过事件驱动实现数据一致性
-
缓存更新
- 实时更新Redis/Elasticsearch缓存
- 避免缓存与数据库数据不一致
-
审计与合规
- 记录所有数据变更历史
- 满足法规要求的数据追踪
环境准备
1 安装 Kafka
# 下载 Kafka wget https://downloads.apache.org/kafka/3.7.0/kafka_2.13-3.7.0.tgz tar -xzf kafka_2.13-3.7.0.tgz cd kafka_2.13-3.7.0 # 启动 Zookeeper bin/zookeeper-server-start.sh config/zookeeper.properties & # 启动 Kafka bin/kafka-server-start.sh config/server.properties &
2 部署 Debezium
# 下载 Debezium Connector wget https://repo1.maven.org/maven2/io/debezium/debezium-connector-mysql/2.5.0.Final/debezium-connector-mysql-2.5.0.Final-plugin.tar.gz tar -xzf debezium-connector-mysql-2.5.0.Final-plugin.tar.gz # 拷贝到 Kafka 插件目录 mkdir /opt/connectors cp -r debezium-connector-mysql/ /opt/connectors/
修改 config/connect-distributed.properties:
bootstrap.servers=localhost:9092 group.id=connect-cluster key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter key.converter.schemas.enable=true value.converter.schemas.enable=true offset.storage.topic=connect-offsets offset.storage.replication.factor=1 config.storage.topic=connect-configs config.storage.replication.factor=1 status.storage.topic=connect-status status.storage.replication.factor=1 plugin.path=/opt/connectors
实战案例:MySQL 数据库变更捕获
1 准备 MySQL 数据库
-- 创建测试数据库
CREATE DATABASE testdb;
-- 创建用户并授权
CREATE USER 'debezium'@'%' IDENTIFIED BY 'password';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'debezium'@'%';
FLUSH PRIVILEGES;
-- 创建测试表
USE testdb;
CREATE TABLE users (
id INT AUTO_INCREMENT PRIMARY KEY,
name VARCHAR(100),
email VARCHAR(100),
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
);
2 注册 Debezium Connector
启动 Kafka Connect:
bin/connect-distributed.sh config/connect-distributed.properties
注册 MySQL Connector:
curl -X POST -H "Content-Type: application/json" http://localhost:8083/connectors -d '{
"name": "mysql-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "localhost",
"database.port": "3306",
"database.user": "debezium",
"database.password": "password",
"database.server.id": "184054",
"database.server.name": "mysql-server-1",
"database.include.list": "testdb",
"table.include.list": "testdb.users",
"database.history.kafka.bootstrap.servers": "localhost:9092",
"database.history.kafka.topic": "dbhistory.users",
"include.schema.changes": "true"
}
}'
3 测试数据捕获
对数据库执行操作:
-- 插入数据
INSERT INTO users (name, email) VALUES ('张三', 'zhangsan@example.com');
INSERT INTO users (name, email) VALUES ('李四', 'lisi@example.com');
-- 更新数据
UPDATE users SET email = 'zhangsan_new@example.com' WHERE id = 1;
-- 删除数据
DELETE FROM users WHERE id = 2;
4 消费 Kafka 消息
# 查看 Topic bin/kafka-topics.sh --list --bootstrap-server localhost:9092 # 输出: # mysql-server-1.testdb.users # 消费消息 bin/kafka-console-consumer.sh \ --bootstrap-server localhost:9092 \ --topic mysql-server-1.testdb.users \ --from-beginning
INSERT 事件示例:
{
"schema": {
"type": "struct",
"fields": [...],
"optional": false,
"name": "mysql-server-1.testdb.users.Envelope"
},
"payload": {
"before": null,
"after": {
"id": 1,
"name": "张三",
"email": "zhangsan@example.com",
"created_at": 1715000000000
},
"source": {
"version": "2.5.0.Final",
"connector": "mysql",
"name": "mysql-server-1",
"ts_ms": 1715000000000,
"snapshot": "false",
"db": "testdb",
"table": "users",
"server_id": 0,
"gtid": null,
"file": "binlog.000003",
"pos": 154,
"row": 0,
"thread": 10,
"query": null
},
"op": "c",
"ts_ms": 1715000000000,
"transaction": null
}
}
UPDATE 事件示例:
{
"payload": {
"before": {
"id": 1,
"name": "张三",
"email": "zhangsan@example.com",
"created_at": 1715000000000
},
"after": {
"id": 1,
"name": "张三",
"email": "zhangsan_new@example.com",
"created_at": 1715000000000
},
"op": "u"
}
}
DELETE 事件示例:
{
"payload": {
"before": {
"id": 2,
"name": "李四",
"email": "lisi@example.com",
"created_at": 1715000001000
},
"after": null,
"op": "d"
}
}
op 字段说明:
c= Create(插入)u= Update(更新)d= Delete(删除)r= Read(快照读取)
进阶案例:同步到 Elasticsearch
1 使用 Kafka Connect Elasticsearch Sink
注册 Elasticsearch Sink Connector:
curl -X POST -H "Content-Type: application/json" http://localhost:8083/connectors -d '{
"name": "elasticsearch-sink",
"config": {
"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
"tasks.max": "1",
"topics": "mysql-server-1.testdb.users",
"connection.url": "http://localhost:9200",
"key.ignore": "false",
"schema.ignore": "true",
"type.name": "_doc",
"behavior.on.malformed.documents": "warn",
"behavior.on.null.values": "delete"
}
}'
2 查询 Elasticsearch
# 查询所有用户
curl -X GET "localhost:9200/mysql-server-1.testdb.users/_search?pretty" -H 'Content-Type: application/json' -d'
{
"query": { "match_all": {} }
}'
使用 Debezium Server 简化部署
对于无需 Kafka 的场景,可以使用 Debezium Server 直接将 CDC 推送到目标系统。
1 配置 Debezium Server 推送到 Redis
# application.properties debezium.sink.type=redis debezium.sink.redis.address=localhost:6379 debezium.sink.redis.password= debezium.sink.redis.database=0 debezium.source.connector.class=io.debezium.connector.mysql.MySqlConnector debezium.source.offset.storage.file.filename=data/offsets.dat debezium.source.database.hostname=localhost debezium.source.database.port=3306 debezium.source.database.user=debezium debezium.source.database.password=password debezium.source.database.server.id=184054 debezium.source.database.server.name=mysql debezium.source.table.include.list=testdb.users debezium.source.database.history.file.filename=data/schema.dat debezium.source.schema.history.internal.file.filename=data/schema-changes.dat
常用操作与管理
1 Connector 管理 API
# 查看所有 Connector curl http://localhost:8083/connectors # 查看 Connector 状态 curl http://localhost:8083/connectors/mysql-connector/status # 暂停 Connector curl -X PUT http://localhost:8083/connectors/mysql-connector/pause # 恢复 Connector curl -X PUT http://localhost:8083/connectors/mysql-connector/resume # 重启 Connector curl -X POST http://localhost:8083/connectors/mysql-connector/restart # 删除 Connector curl -X DELETE http://localhost:8083/connectors/mysql-connector
2 监控指标 (Prometheus Format)
# 查看 Metrics curl -s http://localhost:8083/metrics | grep debezium | head -20
部署架构建议
graph TB
subgraph "生产环境"
DB[(MySQL Master)]
DB2[(MySQL Slave)]
end
subgraph "CDC 集群"
D1[Debezium Connector 1]
D2[Debezium Connector 2]
end
subgraph "消息中间件"
K[Kafka Cluster]
end
subgraph "数据消费端"
E[Elasticsearch]
W[数据仓库]
C[Redis Cache]
S[其他微服务]
end
DB --> D1
DB2 --> D2
D1 --> K
D2 --> K
K --> E
K --> W
K --> C
K --> S
注意事项与最佳实践
- 权限配置:用于片段的数据库用户需要
REPLICATION SLAVE,REPLICATION CLIENT权限 - Binlog 配置:MySQL 需要开启 Binlog,并设置为
ROW格式# my.cnf server-id = 223344 log_bin = mysql-bin binlog_format = ROW binlog_row_image = FULL expire_logs_days = 7
- Schema 变更:Debezium 能够捕获
ALTER TABLE等 DDL 操作 - 大表快照:首次启动时会对现有数据做全量快照,需要预留足够的磁盘空间
- 容错处理:建议为下游消费者配置幂等性处理,以应对重复消息
故障排查
常见问题
| 问题 | 解决方式 |
|---|---|
| Connector 状态为 FAILED | 查看 Connect 日志,检查数据库连接 |
| 收不到数据 | 检查 binlog 是否开启,checkpoint 是否正确 |
| topic 数量过多 | 使用主题路由(Topic Routing)按需配置 |
| 性能瓶颈 | 调整 max.batch.size 和 poll.interval.ms 参数 |
查看日志
# Kafka Connect 日志 tail -f logs/connect.log # 指定 Connector 的日志 grep "mysql-connector" logs/connect.log | tail -50