Debezium 深入
概述
Debezium 是基于 Kafka Connect 的开源 CDC 平台,把数据库变更变成结构化事件流入 Kafka。核心能力:全量快照与增量无缝衔接、多数据库支持、Schema 变更处理。本文讲透 Snapshot 模式、Kafka Connect 集成、Schema 变更与 Exactly-Once。
一、Debezium 定位
| 特性 | 说明 |
|---|---|
| 定位 | 分布式 CDC 平台 |
| 架构 | Kafka Connect Connector |
| 源 | MySQL/PG/Oracle/SQL Server/MongoDB |
| 输出 | Kafka Topic |
| 事件 | 结构化变更(Debezium 格式) |
事件格式:
{
"before": {...},
"after": {...},
"op": "c/u/d/r", -- 操作
"source": {...}, -- 源元数据
"ts_ms": ... -- 时间戳
}二、Kafka Connect 集成
2.1 架构
Kafka Connect(分布式)
├── Source Connector(Debezium 等)
│ 数据库 → Kafka
└── Sink Connector
Kafka → 目标| 组件 | 说明 |
|---|---|
| Connector | 任务定义 |
| Task | 执行单元 |
| Worker | 运行节点 |
| REST API | 管理接口 |
2.2 部署
json
// 提交 Connector
POST /connectors
{
"name": "mysql-orders-connector",
"config": {
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql",
"database.port": "3306",
"database.user": "debezium",
"database.password": "***",
"database.server.name": "orders",
"database.include.list": "orderdb",
"table.include.list": "orderdb.orders",
"database.history.kafka.topic": "dbhistory.orders",
"topic.prefix": "orders"
}
}| 配置 | 说明 |
|---|---|
| connector.class | 连接器类 |
| server.name / topic.prefix | Topic 前缀 |
| include.list | 库表过滤 |
| history.topic | Schema 历史 |
2.3 Topic 结构
topic 命名:{prefix}.{db}.{table}
orders.orderdb.orders
消息 key:主键,value:变更事件三、Snapshot 模式
3.1 全量快照(Initial Snapshot)
首次启动:
1. 记录 binlog 位点(起点)
2. 全量导出当前数据
3. 期间新变更存 binlog
4. 全量完成后从位点续传增量
5. 无缝衔接(快照事件 op=r + 增量)| 阶段 | 说明 |
|---|---|
| 快照 | 全量当前数据 |
| 位点 | 记录起点 |
| 增量 | 续传变更 |
| 衔接 | 无重复无丢失 |
3.2 增量快照(Incremental Snapshot)
大数据量快照分批:
按主键分块逐步抓取
不阻塞写入| 对比 | 全量快照 | 增量快照 |
|---|---|---|
| 数据量 | 一次性 | 分批 |
| 阻塞 | 有影响 | 低 |
| 恢复 | 重来 | 续传 |
3.3 快照一致性
快照与增量使用同一快照位点:
保证一致性快照(MVCC)
不丢不重四、Schema 变更处理
4.1 Schema 历史
数据库 DDL 变更(加列等):
Debezium 记录 Schema 历史
变更事件带新 Schema| 机制 | 说明 |
|---|---|
| Database History | 记录 DDL 历史 |
| Schema 变更事件 | DDL 也发 Kafka |
| 兼容 | 新旧数据混读 |
DDL 处理:
加列 → 事件 schema 更新
下游需适配(如演进)4.2 与 Schema Registry
Avro 格式 + Schema Registry:
统一管理 schema
兼容性检查| 好处 | 说明 |
|---|---|
| 版本管理 | Schema 演进 |
| 兼容检查 | 防破坏 |
| 下游消费 | 自动适配 |
五、Exactly-Once 语义
5.1 挑战
CDC 的 exactly-once:
数据不丢不重
端到端一致5.2 实现层次
| 层 | 保证 |
|---|---|
| Connector 内 | 位点 + 事务 |
| Kafka 内 | 幂等/事务 |
| 端到端 | 幂等下游 |
5.3 关键机制
| 机制 | 说明 |
|---|---|
| Offset | Connector 记录消费位点 |
| 事务 | binlog 事务边界 |
| 幂等 | 事件含主键可去重 |
| 下游 | 事务写入/幂等 |
Exactly-Once 依赖:
Kafka 事务 + 幂等 Sink
或下游唯一键去重5.4 常见语义
| 语义 | 说明 |
|---|---|
| At-least-once | 可能重复(默认) |
| Exactly-once | 需事务配合 |
| 建议 | 幂等兜底 |
六、部署模式
| 模式 | 说明 |
|---|---|
| Single Node | 单机 |
| Distributed | Kafka Connect 集群 |
| Embedded | 内嵌(Flink 等) |
| 场景 | 模式 |
|---|---|
| 简单 | Single |
| 生产 | Distributed |
| Flink | Embedded / Flink CDC |
七、运维要点
| 要点 | 说明 |
|---|---|
| 监控 | 位点/延迟/任务状态 |
| 扩容 | 任务并行 |
| 恢复 | 从位点续传 |
| 多实例 | 表级并发 |
| 幂等 | 下游设计 |
7.1 常见问题
| 问题 | 处理 |
|---|---|
| 延迟大 | 扩 Task/并行 |
| 重复事件 | 下游幂等 |
| DDL 破坏 | Schema 兼容 |
| 快照慢 | 增量快照 |
| 位点丢失 | 重启恢复 |
常见问题速查
| 问题 | 要点 |
|---|---|
| 快照+增量怎么衔接 | 位点记录 |
| 增量快照 | 分批不阻塞 |
| DDL 怎么办 | Schema 历史 |
| Exactly-Once | 事务 + 幂等 |
| 与 Kafka 关系 | Connect 生态 |