Debezium
Debezium 是 Red Hat 发起、现为 Apache 顶级项目 的开源 CDC(Change Data Capture) 平台。它通过解析数据库事务日志(MySQL binlog、PostgreSQL WAL、Oracle redo log 等),将 INSERT / UPDATE / DELETE 变更事件以结构化消息推送到 Kafka,是 Kafka 生态中最常用的 CDC 方案。
架构与组件
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 的对比:
| 维度 | Debezium | Canal |
|---|---|---|
| 数据库 | 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 Registry | Java Client、阿里内部实践多 |
支持的数据库
| Connector | 日志来源 | 典型 Topic 命名 |
|---|---|---|
| MySQL | binlog(ROW) | {prefix}.{db}.{table} |
| PostgreSQL | 逻辑复制 / WAL | {prefix}.{schema}.{table} |
| MongoDB | oplog | {prefix}.{db}.{collection} |
| Oracle | LogMiner / XStream | {prefix}.{schema}.{table} |
| SQL Server | CDC 表 | {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=ONCREATE 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.prefix | Topic 前缀,最终 Topic 为 {prefix}.{db}.{table} |
database.include.list / table.include.list | 白名单过滤 |
snapshot.mode | initial(先全量快照再增量)、never(仅增量)、when_needed |
schema.history.internal.kafka.topic | DDL 历史存储 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
}| 字段 | 含义 |
|---|---|
op | c=INSERT,u=UPDATE,d=DELETE,r=READ(快照阶段) |
before / after | 变更前后行数据;DELETE 时 after 为 null |
source | 源库、表、binlog 位点、时间戳 |
ts_ms | Connector 处理时间(毫秒) |
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-offsetsTopic,重启后从上次位点续读; tasks.max > 1时 MySQL Connector 可按表并行(需配置database.server.id范围)。
监控指标
| 指标 | 说明 |
|---|---|
MilliSecondsBehindSource | Connector 相对源库延迟 |
| 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 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 |
常见使用场景
| 场景 | 说明 |
|---|---|
| 缓存同步 | 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。