CDC(Change Data Capture,变更数据捕获)是把数据库的每一次 insert、update、delete 变成一条可订阅的事件流。它不是「定时查全表对比差异」,而是直接读数据库的预写日志——PostgreSQL 的 WAL、MySQL 的 binlog、MongoDB 的 oplog。Kafka Connect 提供标准化的运行时,Debezium 提供连接器实现,两者组合是当前最主流的开源 CDC 方案。但真正把它跑在生产上,难点从来不在「能不能连上」,而在快照与增量如何无缝衔接、Schema 变更如何不炸下游、以及复制槽丢失后如何恢复。
1. CDC 的定位与日志捕获原理
1.1 为什么从「查表」转向「读日志」
最朴素的同步方式是轮询:每分钟 SELECT * FROM orders WHERE updated_at > :last。它有三个绕不过去的缺陷。
第一,删除无法捕获。物理删除后行消失了,轮询查不到任何痕迹。第二,中间状态丢失。一行在一分钟内被改了五次,轮询只看到最后一次。第三,给源库加压。大表上的增量查询要么走索引(需要 updated_at 索引,且高频更新时索引选择性差),要么全表扫描,都会和生产查询抢资源。
读日志的方式没有这些问题:WAL 里记录了每一次变更的完整前后镜像(取决于 REPLICA IDENTITY),删除是一条 op: "d" 的记录,中间状态一条不落,而且读日志是顺序读,对源库的额外负担极小。代价是 CDC 与具体数据库的日志格式强绑定,运维复杂度上了一个台阶。
1.2 逻辑解码:WAL 到变更事件
以 PostgreSQL 为例。物理复制流传输的是页面的二进制差异,只有同版本的 PostgreSQL 能解读;逻辑复制流传输的是行级逻辑变更,由 pgoutput 之类的逻辑解码插件把 WAL 翻译成可读的变更记录。Debezium 正是通过逻辑复制协议订阅这条流。
要让源库支持逻辑解码,需要满足几个前提:
# postgresql.conf 关键配置
wal_level = logical # 必须为 logical,replica 不够
max_replication_slots = 10 # 至少为 Debezium 连接器数量留够
max_wal_senders = 10
# 需要被捕获的表必须设置 REPLICA IDENTITY
# 默认 default 只记录主键,FULL 记录整行旧值
REPLICA IDENTITY 是个高频踩坑点。默认值 DEFAULT 下,update 与 delete 事件的 before 镜像只包含主键。如果你的下游需要「旧值」,就必须把表改成 REPLICA IDENTITY FULL,代价是 WAL 体积显著增大。这个选择要在写入放大与下游需求之间权衡。
MySQL 侧则依赖 binlog_format = ROW 与 binlog_row_image = FULL,binlog_row_image 默认就是 FULL,所以 MySQL 的 before/after 镜像通常比 PostgreSQL 完整。MongoDB 走 oplog 或 change stream,天然带完整文档。
1.3 初始快照与增量衔接
这是 CDC 最容易出错的地方。一个全新的连接器启动时,历史数据还没进 Kafka,必须先做一次初始快照(snapshot)把存量数据全量导出,然后无缝切换到增量流。Debezium 的默认策略(snapshot.mode = initial)流程如下:
- 在源库上开启一个可重复读(REPEATABLE READ)事务,拿到一个一致性快照点。
- 在快照点位置先创建逻辑复制槽并记录 LSN(Log Sequence Number)。
- 分块扫描所有表,把每行作为
op: "r"(read)事件写入 Kafka。 - 快照完成后,从记录的 LSN 开始消费增量,事件
op变成c/u/d。
关键在于第 2 步与第 3 步的顺序:先占坑,再扫描。这样即使扫描耗时几小时,WAL 也不会被清理掉,因为复制槽会把 WAL 保留住。反过来若先扫描再建槽,扫描期间的变更就丢了。
代价是复制槽会持续保留 WAL,磁盘占用随扫描时间线性增长。大表快照前必须确认磁盘余量,或者改用 snapshot.mode = initial_only 配合后续手工起增量,或者使用增量快照(incremental snapshot,见第 5 节)。
1.4 快照模式的取舍
snapshot.mode 决定连接器启动时对存量数据的处理方式,选错会带来数据丢失或长时间阻塞:
| 模式 | 行为 | 适用场景 |
|---|---|---|
initial | 首次启动做全量快照,之后增量 | 默认,全新接入 |
initial_only | 只做快照,不做增量后停止 | 先导存量、增量另行安排 |
schema_only | 只读表结构,从当前 LSN 起增量 | 下游已有存量数据,只需增量 |
no_data | 同 schema_only | 已废弃别名 |
when_needed | 仅在槽失效或无 offset 时快照 | 自动补救,风险较高 |
never | 从不快照 | 配合外部工具管理位点 |
生产上最容易踩的坑是 schema_only:它跳过快照直接读增量,如果下游其实没有存量数据,就会静默丢历史。反过来 initial 在每次 offset 丢失时都会重跑快照,大表上可能造成数小时的重复投递。建议在 initial 之外显式评估 when_needed,让「槽失效」这一异常路径也有明确的补救动作。
2. Debezium 连接器配置实战
2.1 连接器的核心配置项
以下是一份 PostgreSQL 连接器的生产级配置,关键项都带了注释:
{
"name": "pg-orders-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"tasks.max": "1",
"database.hostname": "pg-primary.internal",
"database.port": "5432",
"database.user": "debezium",
"database.password": "${file:/opt/kafka/secrets.properties:pg_password}",
"database.dbname": "shop",
"topic.prefix": "cdc",
"table.include.list": "public.orders,public.order_items",
"plugin.name": "pgoutput",
"slot.name": "debezium_orders",
"publication.autocreate.mode": "filtered",
"snapshot.mode": "initial",
"heartbeat.interval.ms": "10000",
"decimal.handling.mode": "string",
"time.precision.mode": "connect",
"tombstones.on.delete": "true",
"transforms": "unwrap",
"transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
"transforms.unwrap.drop.tombstones": "false"
}
}
几个必须理解的点:
tasks.max对 Debezium 基本无意义。CDC 连接器只能单任务运行,因为复制槽的位点是单一的。设成大于 1 会启动失败或产生多个重复槽。topic.prefix决定 topic 命名前缀,实际 topic 形如cdc.public.orders。slot.name必须唯一。多个连接器共用一个槽会互相抢位点,导致数据错乱。heartbeat.interval.ms是防复制槽膨胀的关键,详见第 5 节。ExtractNewRecordState(unwrap) 把 Debezium 的复杂信封压平成业务字段,让下游消费者不用解析before/after。但它丢弃了 op 类型与旧值,需要审计场景时不能开。
2.2 事件信封结构
不开启 unwrap 时,每条消息的 value 是一层信封:
{
"before": { "id": 1001, "status": "created", "amount": "99.00" },
"after": { "id": 1001, "status": "paid", "amount": "99.00" },
"source": {
"version": "2.5.0.Final",
"connector": "postgresql",
"db": "shop", "schema": "public", "table": "orders",
"lsn": 24023128, "txId": 1876
},
"op": "u",
"ts_ms": 1759800000000
}
op 的取值是 CDC 语义的核心:c 创建、u 更新、d 删除、r 快照读。source.lsn 提供了位点信息,下游若要做精确去重,可以用 (source.lsn, source.txId) 作为幂等键。before 与 after 同时存在时,diff 出真正变化的字段能减少下游无效更新。
2.3 事务边界与消息顺序
Debezium 默认在事务提交后才把该事务内的事件批量投递,因此同一个事务的变更在同一个分区内保持原序,且不会出现「半个事务」。这依赖 topic.transaction 与连接器内部的缓冲。
需要注意的是:顺序保证仅在单分区内成立。同一个表的所有事件会路由到同一个 topic,但若下游按 key 分区(比如按 order_id),同一行的先后顺序仍然保持,跨行的顺序则不再保证。设计下游逻辑时不要依赖跨行顺序。
2.4 MySQL 与 MongoDB 的配置差异
换数据库时连接器类名与关键参数都要改,但语义上的差异更值得注意:
{
"connector.class": "io.debezium.connector.mysql.MySqlConnector",
"database.hostname": "mysql-primary.internal",
"database.server.id": "184054",
"database.include.list": "shop",
"table.include.list": "shop.orders",
"schema.history.internal.kafka.topic": "schema-changes.shop",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"snapshot.locking.mode": "minimal"
}
MySQL 侧有三个特有概念:
database.server.id必须在整个复制拓扑中唯一,重复会直接踢掉其他副本。schema.history.internal.kafka.topic用来存 DDL 历史。MySQL 的 binlog 记录 DDL 但不含表结构快照,Debezium 必须自己维护一份 DDL 历史才能在重启后还原 Schema,这个 topic 丢了就得重做快照。snapshot.locking.mode控制快照期间的加锁策略。默认minimal只在快照开始瞬间持全局读锁,none完全不锁(可能读到不一致数据),extended全程持锁(阻塞写入)。
MongoDB 则用 MongoDbConnector,走 change stream(4.0+)或 oplog,事件结构里 before/after 是完整文档,且没有 Schema 注册表可依,通常直接用 JSON converter。跨库复制的常见做法是把 PostgreSQL 的 CDC 流投递到 Apache Flink 流处理
做转换后再写目标库。
3. Schema 演进与兼容性
3.1 为什么 CDC 必须配 Schema Registry
数据库表的列会变。ALTER TABLE orders ADD COLUMN coupon_id bigint 之后,Debezium 立刻开始在新事件里带上 coupon_id。如果下游消费者用的是强类型反序列化(Avro、Protobuf),Schema 不匹配会直接抛异常,消费停摆。
Schema Registry 承担三件事:集中存储每个 topic 的 Schema 版本、在生产者注册时校验兼容性、给每条消息附一个 schema id(而不是把完整 Schema 塞进消息,节省带宽)。Debezium 原生支持 Avro 与 JSON Schema 两种 converter,通过 key.converter / value.converter 配置。
{
"key.converter": "io.confluent.connect.avro.AvroConverter",
"key.converter.schema.registry.url": "http://schema-registry:8081",
"value.converter": "io.confluent.connect.avro.AvroConverter",
"value.converter.schema.registry.url": "http://schema-registry:8081",
"value.converter.schemas.enable": "true"
}
关于 Schema Registry 本身的注册、兼容性校验与 REST API,可以参考 Kafka Schema Registry 与 Avro 实践 。
3.2 兼容性策略怎么选
Schema Registry 支持四种兼容性级别,选错会直接卡住 DDL 发布:
| 级别 | 允许的变更 | 适用场景 |
|---|---|---|
| BACKWARD | 删字段、加可选字段 | 消费者先升级(默认值最宽松) |
| FORWARD | 加字段、删可选字段 | 生产者先升级 |
| FULL | 两者的交集 | 双方独立升级,最安全 |
| NONE | 任意 | 仅调试,生产禁用 |
CDC 场景下推荐 BACKWARD 起步、逐步收紧到 FULL。因为 Debezium 是「生产者」,而生产者的升级节奏由 DDL 决定,往往不可控;消费者则需要时间适配。BACKWARD 允许「加可选字段」,正好覆盖 ADD COLUMN 这个最高频的变更。
必须警惕的是 DROP COLUMN 与类型变更。DROP COLUMN 在 BACKWARD 下是允许的(删字段),但会破坏那些仍在读该字段的老消费者;类型变更(int 改 bigint、varchar 改 text)几乎总是破坏性变更,会被直接拒绝。此时唯一安全的做法是新建一个字段,做双写迁移,等所有消费者切换后再删旧字段。
3.3 DDL 变更的连锁反应
PostgreSQL 上还有一个隐蔽问题:Debezium 捕获 DDL 本身的能力有限。它不会把 ALTER TABLE 作为事件发出去,而是在下一次数据变更时以新 Schema 发送。这意味着下游无法通过事件流感知表结构变化,只能依赖 Schema Registry 的版本号。
更麻烦的是列重命名。Debezium 看到的是「旧列消失、新列出现」,而 Schema Registry 看到的是「删一列、加一列」。如果兼容性设为 FULL,这次变更会被拒绝,连接器报错停摆。生产上要么提前把级别放宽,要么避免直接 rename,改用「加新列 + 数据回填 + 删旧列」的三步走。
4. 重复、顺序与精确一次
4.1 at-least-once 的来源
Debezium 默认语义是 at-least-once,重复主要来自两个环节:
其一,快照与增量的边界。 连接器在快照完成后从记录的 LSN 开始读增量。如果连接器在快照期间崩溃重启,会从头重跑快照(取决于 snapshot.mode),已写入的事件会重复。
其二,offset 提交与事件投递的时序。 连接器先把事件写进 Kafka,再提交 source offset。若在写完之后、提交之前崩溃,重启后会从上一个已提交位点重读,导致一批事件重复投递。
这两处的重复都是结构性的,无法通过配置消除,只能在消费端做幂等。
4.2 下游幂等设计
幂等键的选取决定去重是否可靠。推荐用 (source.lsn, source.txId, source.table) 组合,或者业务主键 + 版本号。以下是一个 Flink 侧按主键 upsert 的写法,天然幂等:
// 把 CDC 流按主键 upsert 到外部存储,天然去重
KafkaSource<String> source = KafkaSource.<String>builder()
.setBootstrapServers("kafka:9092")
.setTopics("cdc.public.orders")
.setGroupId("cdc-orders-sink")
.setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
// 关键:以 order_id 为主键,重复事件覆盖同一行,结果一致
stream.keyBy(e -> e.orderId).process(new UpsertSink());
若下游是纯 append 的日志(比如追加到 Hive 或对象存储),则必须显式去重:维护一个「已处理主键 + 版本」的状态表,遇到旧版本直接丢弃。这需要额外的存储,且状态会随数据量增长。
4.3 精确一次的两个层次
有人会问:Kafka Connect 不是支持 exactly-once 吗?是的,但只覆盖 Kafka 内部的读写。开启 exactly.once.support = required 后,连接器用 Kafka 事务保证「事件写入 + offset 提交」原子化,消除了 4.1 中第二类重复。但第一类(快照重跑)以及「写 Kafka 到写外部系统」这一段仍然不在事务内。
因此「端到端精确一次」的实际含义是:Kafka 内部精确一次 + 外部系统幂等写入。任何声称「开了 exactly-once 就万事大吉」的说法都忽略了后半段。关于事务与幂等生产者的底层机制,可以看 Kafka 事务与 Exactly-Once 语义 。
5. 断点续传与故障恢复
5.1 offset 的存储与重启行为
Debezium 的 source offset 存在 Kafka 的内部 topic(connect-offsets)里,内容是逻辑复制槽名与 LSN。连接器重启时:
- 读取已提交的 offset。
- 检查复制槽是否存在。若存在,从该槽的
confirmed_flush_lsn继续。 - 若槽不存在(被手工删除或数据库重建),则根据
snapshot.mode决定:initial会重跑快照,schema_only会跳过快照直接尝试从当前 LSN 开始(可能丢数据)。
这里的陷阱是:连接器 offset 与复制槽是两套状态,可能不一致。比如连接器重建但复制槽还在,Debezium 会优先信任槽的位点;反之槽丢了但 offset 还在,会触发快照重跑。运维时必须同时检查两者。
5.2 复制槽膨胀与丢失
膨胀(bloat) 是最常见的生产事故。复制槽会保留所有尚未被消费的 WAL。若连接器长时间下线(网络故障、Kafka 不可用),WAL 会持续堆积,直到把源库磁盘写满,导致源库写入全停。这是 CDC 最危险的一类故障,因为它会拖垮生产数据库。
防护手段有三个层次:
# 1. 心跳:让连接器在无数据变更时也推进槽位点
heartbeat.interval.ms = 10000
# 2. 数据库侧兜底:PostgreSQL 13+ 支持
max_slot_wal_keep_size = 10GB
# 超过上限时槽会被标记为不可用,WAL 得以回收,
# 但连接器恢复时必须重做快照(数据不丢,只是要重跑)
# 3. 监控:把槽的滞后量做成告警
SELECT slot_name,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS lag
FROM pg_replication_slots;
max_slot_wal_keep_size 是一个「两害相权」的开关:不设,源库有被写满的风险;设了,槽可能失效需要重做快照。生产上通常设置一个能容忍的阈值(如 10~20 GB),并配监控告警。
槽丢失 的典型原因:DBA 手工 pg_drop_replication_slot 清理、主从切换后槽未同步到新主、或者 max_slot_wal_keep_size 触发。恢复路径是删除连接器的 offset 并重建,触发一次全量快照——前提是业务能接受重跑窗口。
5.3 增量快照:不停机的一致性补救
传统快照会锁表(PostgreSQL 上虽用可重复读而非锁表,但仍会长时间占用一个事务、产生大量 WAL),对超大表不友好。Debezium 1.6 起提供的增量快照(incremental snapshot,KIP-650 相关信号机制)解决了这个问题:
它按主键分块(chunk)读取,每块读完后短暂停顿,让增量事件穿插进来,因此快照期间不阻塞、可中断、可恢复。触发方式是通过 debezium-signal topic 发送一个 signal:
{
"id": "ad-hoc-1",
"type": "execute-snapshot",
"data": {
"data-collections": ["public.orders"],
"type": "incremental",
"additional-condition": "status = 'open'"
}
}
增量快照读到的行会带 op: "r",与增量事件按 LSN 去重后合流。它是「事后补救」——比如某张表漏配了 table.include.list、或者数据不一致需要重刷——最优雅的方案,避免了「停连接器 + 删 offset + 全量重跑」的重操作。
6. 监控与排错
6.1 必看的指标
CDC 的健康度不能只看「连接器是不是 RUNNING」。核心指标分三类:
- 连接器侧:
MilliSecondsBehindSource(落后源库多久)、NumberOfEventsFiltered、TotalNumberOfEventsSeen、任务状态与重启次数。 - 复制槽侧:
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)的字节滞后,以及active标志。 - Kafka 侧:CDC topic 的 consumer lag、分区是否倾斜。
MilliSecondsBehindSource 持续增长往往意味着某个 topic 写入被限流(配额或 broker 过载),或者反序列化失败在重试。关于消费者滞后的一般排查方法,可参考 Kafka 消费延迟诊断
。
6.2 典型故障与处置
故障一:连接器反复重启,日志报 replication slot already exists。 原因是上一个连接器没清理干净。处置:确认旧连接器已删除、旧槽无进程占用后,SELECT pg_drop_replication_slot('debezium_orders') 再重建。
故障二:Schema 不兼容导致任务 FAILED。 日志里会有 Schema being registered is incompatible with an earlier schema。处置:临时把该 subject 的兼容级别调成 NONE 让变更通过(有风险),或按第 3.3 节的三步走改造 DDL。
故障三:数据重复。 先确认是不是快照重跑导致的整体重复(可按 op: "r" 过滤),若是增量重复则检查下游幂等键。不要试图在 Kafka 侧「去重」,那是治标。
故障四:槽膨胀触发源库磁盘告警。 立即恢复连接器消费,或提高 max_slot_wal_keep_size 让 WAL 回收,同时准备重做快照。
6.3 吞吐与资源调优
CDC 连接器的吞吐瓶颈通常不在 Kafka,而在源库的 WAL 生成速率与连接器的单任务处理能力。可调的杠杆有限:
max.batch.size与max.queue.size(默认 2048 / 8192)。增大能提升单次批量,但会占用更多堆内存;队列满时连接器会阻塞读取,反而拖慢复制槽推进。poll.interval.ms(默认 500)。缩短能更快感知变更,代价是空轮询变多。table.include.list收窄。捕获的表越少,WAL 解析与事件构造的开销越小。decimal.handling.mode用string而非double,避免精度丢失,也省去 BigDecimal 序列化开销。
一个常见的误解是「加 Kafka 分区就能提速」。CDC 连接器是单任务,写入的分区数由表的 key 决定,加分区不会提升连接器吞吐,只会改变下游并行度。真正的提速手段是拆分连接器:按业务域把不同的表分给不同的连接器(各自独立的复制槽),水平扩展捕获能力。
此外,ExtractNewRecordState 这类 SMT 会逐条做转换,在大流量下是实打实的 CPU 开销。若下游能接受信封结构,关掉 SMT 能省下可观的 CPU。
7. 常见坑清单
- 忘记设置
wal_level = logical,连接器启动直接失败。 REPLICA IDENTITY保持默认,下游拿不到 update/delete 的旧值。- 多个连接器共用
slot.name,位点互相覆盖导致数据错乱。 tasks.max > 1,Debezium 单任务约束被违反。- 开 unwrap 后误以为还能拿到
before镜像与op。 - 兼容性级别设为
NONE忘了改回来,后续破坏性变更畅通无阻。 - 没有配
heartbeat.interval.ms,空闲表导致槽位点不推进、WAL 堆积。 - 大表快照前没估磁盘,快照跑一半源库写满。
- 主从切换后新主没有对应复制槽,连接器恢复时悄悄重跑快照。
- 依赖跨行顺序:CDC 只保证单分区内单表的事务顺序。
8. 总结
Debezium + Kafka Connect 把 CDC 的「读日志」部分做得很扎实,但生产可用性取决于你如何处理它没有替你解决的部分:快照与增量的衔接顺序、复制槽的容量风险、Schema 演进的下游兼容、以及 at-least-once 的幂等兜底。
工程上建议按这个顺序落地:先确认源库日志参数与 REPLICA IDENTITY,再配好 Schema Registry 与兼容性级别,接着把复制槽滞后纳入监控告警,最后在消费端实现幂等写入。这四步都做完,CDC 才算真正可靠。若下游还要做流式聚合与 join,可以把变更流直接喂给 Kafka Streams 流处理
,省去中间落库的往返。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。