Kafka Connect 与 Debezium CDC 实战

以 Debezium 与 Kafka Connect 为核心讲透 CDC 生产落地:逻辑解码与初始快照衔接、Schema Registry 兼容性策略、重复与顺序保证、复制槽膨胀与断点续传、增量快照补救、监控指标与典型故障处置、下游幂等与吞吐调优

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)流程如下:

  1. 在源库上开启一个可重复读(REPEATABLE READ)事务,拿到一个一致性快照点。
  2. 在快照点位置先创建逻辑复制槽并记录 LSN(Log Sequence Number)。
  3. 分块扫描所有表,把每行作为 op: "r"(read)事件写入 Kafka。
  4. 快照完成后,从记录的 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。连接器重启时:

  1. 读取已提交的 offset。
  2. 检查复制槽是否存在。若存在,从该槽的 confirmed_flush_lsn 继续。
  3. 若槽不存在(被手工删除或数据库重建),则根据 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 流处理 ,省去中间落库的往返。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 集群升级与滚动重启实践
  2. Kafka 应用测试策略:Testcontainers 与集成测试
  3. 压缩算法选型:lz4、zstd、snappy 与 gzip