Debezium案例

wen java案例 4

Debezium 案例详解

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

Debezium案例


核心概念

graph LR
    A[数据库 MySQL/PostgreSQL/MongoDB] -->|捕获变更| B[Debezium Connector]
    B -->|发送到| C[Kafka Topic]
    C -->|消费| D[下游应用/数据仓库/搜索引擎]

典型应用场景

  1. 数据库复制与迁移

    • 从旧数据库实时复制到新数据库
    • 构建数据湖/数据仓库的实时ETL
  2. 微服务数据同步

    • 多个微服务共享数据但各自维护独立的数据库
    • 通过事件驱动实现数据一致性
  3. 缓存更新

    • 实时更新Redis/Elasticsearch缓存
    • 避免缓存与数据库数据不一致
  4. 审计与合规

    • 记录所有数据变更历史
    • 满足法规要求的数据追踪

环境准备

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

注意事项与最佳实践

  1. 权限配置:用于片段的数据库用户需要 REPLICATION SLAVE, REPLICATION CLIENT 权限
  2. Binlog 配置:MySQL 需要开启 Binlog,并设置为 ROW 格式
    # my.cnf
    server-id = 223344
    log_bin = mysql-bin
    binlog_format = ROW
    binlog_row_image = FULL
    expire_logs_days = 7
  3. Schema 变更:Debezium 能够捕获 ALTER TABLE 等 DDL 操作
  4. 大表快照:首次启动时会对现有数据做全量快照,需要预留足够的磁盘空间
  5. 容错处理:建议为下游消费者配置幂等性处理,以应对重复消息

故障排查

常见问题

问题 解决方式
Connector 状态为 FAILED 查看 Connect 日志,检查数据库连接
收不到数据 检查 binlog 是否开启,checkpoint 是否正确
topic 数量过多 使用主题路由(Topic Routing)按需配置
性能瓶颈 调整 max.batch.sizepoll.interval.ms 参数

查看日志

# Kafka Connect 日志
tail -f logs/connect.log
# 指定 Connector 的日志
grep "mysql-connector" logs/connect.log | tail -50

抱歉,评论功能暂时关闭!