CDC 变更数据捕获:Debezium+Kafka 同步 MySQL 数据到数据湖

1. 核心概念解析
  • CDC (变更数据捕获)
    实时捕获数据库变更(增删改),避免全量扫描。MySQL 通过 binlog 实现,捕获效率为 $O(1)$。
  • Debezium
    开源 CDC 工具,将 binlog 转为事件流,支持 Exactly-Once 语义。
  • 数据湖
    集中式存储库(如 S3/HDFS),支持结构化/半结构化数据。
2. 架构流程
graph LR
  MySQL -->|binlog| Debezium 
  Debezium -->|Avro/JSON| Kafka 
  Kafka -->|Connector| 数据湖(S3/Hudi/Iceberg)

3. 实现步骤
(1) 环境配置
  • MySQL 启用 binlog
    配置 my.cnf
    [mysqld]
    server-id=1
    log_bin=mysql-bin
    binlog_format=ROW
    

  • Debezium 连接器部署
    通过 Kafka Connect 注册:
    {
      "name": "mysql-connector",
      "config": {
        "connector.class": "io.debezium.connector.mysql.MySqlConnector",
        "database.hostname": "mysql-host",
        "database.user": "debezium",
        "database.password": "密码",
        "database.server.id": "184054",
        "database.server.name": "dbserver1",
        "table.include.list": "public.*"
      }
    }
    

(2) 数据流处理
  • Debezium 输出结构
    {
      "before": {...},  // 变更前数据
      "after": {...},   // 变更后数据
      "op": "c/u/d"     // 操作类型(增/改/删)
    }
    

  • Kafka Topic 分区策略
    按表主键哈希分区,保证事件顺序性。
(3) 写入数据湖

使用 S3 Sink Connector 示例配置:

{
  "name": "s3-sink",
  "config": {
    "connector.class": "io.confluent.connect.s3.S3SinkConnector",
    "s3.bucket.name": "my-data-lake",
    "storage.class": "io.confluent.connect.s3.storage.S3Storage",
    "format.class": "io.confluent.connect.s3.format.parquet.ParquetFormat",
    "partitioner.class": "io.confluent.connect.storage.partitioner.HourlyPartitioner",
    "flush.size": "10000"
  }
}

4. 关键优化策略
  • 数据一致性
    通过 Kafka 的 ISR 机制 保证,满足 $R + W > N$(N=副本数,R=读副本,W=写确认数)。
  • Schema 演进
    使用 Avro + Schema Registry 管理表结构变更。
  • 压缩存储
    数据湖采用 Parquet 列式存储,空间节省比达 $1:5 \sim 1:10$。
5. 故障处理
  • 断点续传
    Debezium 持久化 offset,Kafka 消费位点自动恢复。
  • 死信队列
    配置 errors.tolerance=all 将异常数据写入特定 Topic。

:实际部署需考虑网络延迟、数据湖分区策略(如按日期分桶)、及 GDPR 数据脱敏要求。

更多推荐