Debezium

Debezium 是 Red Hat 发起、现为 Apache 顶级项目 的开源 CDC(Change Data Capture) 平台。它通过解析数据库事务日志(MySQL binlog、PostgreSQL WAL、Oracle redo log 等),将 INSERT / UPDATE / DELETE 变更事件以结构化消息推送到 Kafka,是 Kafka 生态中最常用的 CDC 方案。

→ CDC 总览 · Canal(MySQL 专用替代)


架构与组件

MySQL / PostgreSQL / MongoDB / ...
     │  binlog / WAL / oplog
     ▼
Debezium Connector(Kafka Connect Source Plugin)
     │  解析日志 → Change Event
     ▼
Kafka Connect Worker(分布式或单机)
     │  写入 Topic + Schema History Topic
     ▼
Kafka Topic(如 myapp.inventory.users)
     │
     ├── Spring @KafkaListener → Redis / ES
     ├── Flink / Kafka Streams → 实时 ETL
     └── Sink Connector → JDBC / Elasticsearch
组件说明
Debezium Connector针对各数据库的 Source 插件(MySQL、PostgreSQL、MongoDB 等)
Kafka Connect运行 Connector 的框架,负责配置分发、任务调度、offset 持久化
Schema History Topic存储 DDL 历史,Connector 重启后恢复表结构
Change Data Topic每张表(或整库)对应一个 Topic,消息为 Debezium Envelope
Schema Registry(可选)配合 Avro / Protobuf 序列化,管理 Schema 版本演化

与 Canal 的对比:

维度DebeziumCanal
数据库MySQL / PG / Oracle / MongoDB / SQL Server 等仅 MySQL
输出Kafka(Kafka Connect 标准)Kafka / RocketMQ / TCP 直连
部署依赖 Kafka Connect 集群独立 Canal Server
消息格式统一 Envelope(before/after/op/source)自定义 JSON / Protobuf
生态Kafka Connect Sink、Flink、Schema RegistryJava Client、阿里内部实践多

支持的数据库

Connector日志来源典型 Topic 命名
MySQLbinlog(ROW){prefix}.{db}.{table}
PostgreSQL逻辑复制 / WAL{prefix}.{schema}.{table}
MongoDBoplog{prefix}.{db}.{collection}
OracleLogMiner / XStream{prefix}.{schema}.{table}
SQL ServerCDC 表{prefix}.{schema}.{table}

快速上手(MySQL → Kafka)

1. MySQL 前提

[mysqld]
server-id=1
log-bin=mysql-bin
binlog-format=ROW
binlog-row-image=FULL
expire_logs_days=7          # 至少覆盖 Connector 最长停机时间
gtid_mode=ON                # 推荐,便于主从切换后续传
enforce_gtid_consistency=ON
CREATE USER 'debezium'@'%' IDENTIFIED BY 'dbz';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT
  ON *.* TO 'debezium'@'%';
FLUSH PRIVILEGES;

2. 启动 Kafka Connect(Docker Compose 示例)

services:
  connect:
    image: quay.io/debezium/connect:2.7
    ports:
      - "8083:8083"
    environment:
      BOOTSTRAP_SERVERS: kafka:9092
      GROUP_ID: connect-cluster
      CONFIG_STORAGE_TOPIC: connect-configs
      OFFSET_STORAGE_TOPIC: connect-offsets
      STATUS_STORAGE_TOPIC: connect-status
      KEY_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      VALUE_CONVERTER: org.apache.kafka.connect.json.JsonConverter
      VALUE_CONVERTER_SCHEMAS_ENABLE: "false"

3. 注册 MySQL Connector

curl -X POST http://localhost:8083/connectors \
  -H "Content-Type: application/json" \
  -d '{
    "name": "inventory-connector",
    "config": {
      "connector.class": "io.debezium.connector.mysql.MySqlConnector",
      "tasks.max": "1",
      "database.hostname": "mysql",
      "database.port": "3306",
      "database.user": "debezium",
      "database.password": "dbz",
      "database.server.id": "184054",
      "topic.prefix": "myapp",
      "database.include.list": "inventory",
      "table.include.list": "inventory.users,inventory.orders",
      "schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
      "schema.history.internal.kafka.topic": "schemahistory.myapp",
      "snapshot.mode": "initial",
      "include.schema.changes": "true",
      "tombstones.on.delete": "true"
    }
  }'

常用配置项:

配置说明
topic.prefixTopic 前缀,最终 Topic 为 {prefix}.{db}.{table}
database.include.list / table.include.list白名单过滤
snapshot.modeinitial(先全量快照再增量)、never(仅增量)、when_needed
schema.history.internal.kafka.topicDDL 历史存储 Topic,不可删
database.server.id模拟 Slave 的 server-id,不能与集群冲突
heartbeat.interval.ms无变更时也发心跳,便于 lag 监控

4. 查看 Connector 状态

curl http://localhost:8083/connectors/inventory-connector/status
curl http://localhost:8083/connectors/inventory-connector/config

消息格式(Debezium Envelope)

Kafka 消息 key 通常为主键 JSON;value 为 Envelope:

{
  "before": {
    "id": 1,
    "name": "Alice",
    "email": "alice@old.com"
  },
  "after": {
    "id": 1,
    "name": "Alice",
    "email": "alice@new.com"
  },
  "source": {
    "version": "2.7.0.Final",
    "connector": "mysql",
    "name": "myapp",
    "ts_ms": 1716000000000,
    "snapshot": "false",
    "db": "inventory",
    "table": "users",
    "server_id": 1,
    "file": "mysql-bin.000123",
    "pos": 456789,
    "row": 0,
    "thread": 12,
    "query": null
  },
  "op": "u",
  "ts_ms": 1716000000100,
  "transaction": null
}
字段含义
opc=INSERT,u=UPDATE,d=DELETE,r=READ(快照阶段)
before / after变更前后行数据;DELETE 时 after 为 null
source源库、表、binlog 位点、时间戳
ts_msConnector 处理时间(毫秒)
transaction同一事务内多条变更共享 transaction id

DELETE + tombstone:tombstones.on.delete=true 时,DELETE 事件后 Kafka 会写入同 key 的 tombstone,Compact Topic 可保留最新状态。


PostgreSQL Connector 要点

PostgreSQL 使用逻辑复制而非直接读 binlog:

-- postgresql.conf
wal_level = logical
max_replication_slots = 4
 
-- 创建专用用户
CREATE USER debezium WITH REPLICATION LOGIN PASSWORD 'dbz';
GRANT SELECT ON ALL TABLES IN SCHEMA public TO debezium;
 
-- 创建 publication(PG 10+)
CREATE PUBLICATION dbz_pub FOR TABLE users, orders;

Connector 配置差异:

{
  "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
  "database.hostname": "postgres",
  "database.port": "5432",
  "database.user": "debezium",
  "database.password": "dbz",
  "database.dbname": "shop",
  "topic.prefix": "myapp",
  "plugin.name": "pgoutput",
  "publication.name": "dbz_pub",
  "slot.name": "debezium_slot"
}

逻辑复制槽(replication slot)会阻止 WAL 被回收;Connector 长时间离线时需监控 slot lag,避免磁盘撑满。


快照模式(Snapshot Mode)

首次启动或 snapshot.mode=initial 时,Connector 会对源表做一致性快照(类似 mysqldump + binlog 位点锁定):

模式行为
initial先全量快照,完成后切增量(默认,适合新链路)
initial_only只做快照,不读增量
never仅增量,要求已有可靠起始位点
when_needed无 offset 时自动快照
no_data快照结构但不导数据

生产注意:

  • 快照期间对源库有 SELECT 读负载,大表可配置 snapshot.fetch.size 分批读取;
  • 快照与增量切换点由 binlog 位点保证,下游需能处理 op=r 的 READ 事件或过滤掉;
  • 断连后若 binlog 已被清理,需重新 initial 或手动指定 snapshot.mode=schema_only + 外部全量。

Schema 演化与 DDL

Debezium 将 DDL 写入 Schema History Topic,并在变更 Topic 中可选附带 schema 信息:

  • include.schema.changes=true:DDL 作为独立事件写入 {prefix} Topic;
  • 配合 Confluent Schema Registry + Avro:字段增删改由 Registry 管理版本,Sink 端自动兼容;
  • ALTER TABLE 可能导致 Connector 短暂重启任务;复杂 DDL(改主键、拆表)需提前演练。

Spring Boot 消费示例

Topic 名 myapp.inventory.users,按主键保证分区顺序:

@KafkaListener(
    topics = "myapp.inventory.users",
    groupId = "cache-sync",
    containerFactory = "manualAckContainerFactory"
)
public void onUserChange(ConsumerRecord<String, String> record, Acknowledgment ack) {
    JsonNode envelope = objectMapper.readTree(record.value());
    String op = envelope.path("op").asText();
 
    if ("d".equals(op)) {
        JsonNode before = envelope.path("before");
        redisTemplate.delete("user:" + before.path("id").asLong());
    } else {
        JsonNode after = envelope.path("after");
        // INSERT / UPDATE:删缓存或按 version 更新
        redisTemplate.delete("user:" + after.path("id").asLong());
    }
 
    ack.acknowledge();  // 先处理成功再提交 offset
}

幂等与顺序要点见 数据同步与 CDC:

  • 以 主键 作为 Kafka message key,保证同行变更顺序;
  • 用 source.file + source.pos + source.row 或 GTID 构造 eventId 去重;
  • 更新类操作比较业务 version,拒绝旧事件覆盖新状态。

高可用与运维

Kafka Connect 集群

Kafka Connect Worker 1 ─┐
Kafka Connect Worker 2 ─┼─ connect-configs / connect-offsets / connect-status
Kafka Connect Worker 3 ─┘
  • Connector 任务可在 Worker 间 rebalance;
  • offset 持久化在 connect-offsets Topic,重启后从上次位点续读;
  • tasks.max > 1 时 MySQL Connector 可按表并行(需配置 database.server.id 范围)。

监控指标

指标说明
MilliSecondsBehindSourceConnector 相对源库延迟
Kafka consumer lag各 Change Topic 堆积
MySQL SHOW MASTER STATUS / PG replication slot lag源端日志是否即将被清理
Snapshot 进度JMX / Connect REST status

常见问题

问题处理
binlog 已被 purge重新全量快照或从备份恢复位点
Schema History Topic 被误删Connector 无法恢复 DDL,需重建 Connector
大事务 DELETE 百万行背压、增大 Connect 内存、下游批量处理
主从切换启用 GTID;确认新主 binlog 连续
重复消息消费端幂等(at-least-once 语义)

Flink CDC(flink-connector-mysql-cdc)底层同样基于 Debezium 引擎,但嵌入 Flink 作业,无需 Kafka 中转:

MySQL ──► Flink CDC Source ──► Flink 算子 ──► Kafka / JDBC / ES
选型适合
Debezium + Kafka Connect多下游订阅、与 Kafka 生态深度集成、运维 Connect 集群
Flink CDC全程流式 ETL、有状态聚合、无需独立 Connect

详见 Flink、CDC。


常见使用场景

场景说明
缓存同步DB 变更 → 删 Redis key,见 缓存与一致性
搜索同步MySQL / PG → Kafka → Elasticsearch Sink
异构同步MySQL → ClickHouse / MongoDB(Sink Connector 或自研消费)
微服务读模型Transactional Outbox 或业务表 CDC → CQRS 读库
数据迁移全量快照 + 增量追平,停机窗口趋近于零
实时数仓 ODS替代 T+1 批同步,见 数仓架构

注意事项

  • binlog / WAL 必须是 ROW / logical 级别,否则拿不到行级前后镜像;
  • database.server.id(MySQL)不能与现有 Slave 冲突;
  • Schema History Topic 和 connect-offsets 禁止随意删除;
  • 初次 initial 快照对大表有读压力,可业务低峰注册 Connector;
  • DDL 变更需消费端兼容或配合 Schema Registry;
  • 端到端按 at-least-once + 幂等 设计,不要假设 exactly-once。

相关链接

同目录

中间件

  • Kafka — Connect 运行载体
  • Flink — Flink CDC 基于 Debezium 引擎

Spring Cloud 实战

数据库