Debezium PostgreSQL 连接器 — CDC 实战
内容提要
Debezium 是分布式 CDC 平台,通过追踪 PostgreSQL 的 WAL 日志,将行变更实时发布到 Kafka。本文介绍如何用 Docker 搭建 PostgreSQL、Kafka 和 Kafka Connect 环境,注册 Debezium 连接器,并演示初始快照和实时增删改事件。事件包含 before、after、source 和 op 字段,主键作为消息键保证顺序。Debezium 读取 WAL 而非表,对应用查询无影响,仅增加 WAL 保留开销。
延伸解读
WAL 逻辑解码与复制槽
Debezium 依赖 PostgreSQL 的逻辑解码功能,通过创建复制槽(replication slot)来追踪 WAL 的读取位置。复制槽会阻止 PostgreSQL 清理尚未被消费的 WAL 段,因此如果 Debezium 长时间停止或故障,WAL 可能积压,占用磁盘空间。配置中需设置 wal_level=logical,并注意 max_replication_slots 和 max_wal_senders 参数,确保有足够的槽位和发送进程。
事件结构与主键分区
每个 Debezium 事件都包含 before、after、source 和 op 字段,op 表示操作类型(r/c/u/d)。默认情况下,UPDATE 和 DELETE 事件的 before 为 null,因为 PostgreSQL 的默认副本标识(replica identity)只记录主键,若需完整旧值,需将表副本标识设为 FULL。消息键使用主键,保证同一行的变更事件进入同一 Kafka 分区,从而维持顺序。
快照与流式切换
连接器首次启动时,snapshot.mode=initial 会先执行全量快照,将现有行以 op:r 事件发出,然后切换到实时流式读取。快照期间产生的变更会通过 WAL 在快照后继续捕获,确保不丢失。快照事件与流式事件在 source 字段中通过 snapshot 值区分(如 first 或 false),便于下游识别。
对应用查询无影响
Debezium 通过复制连接读取 WAL,而非直接查询表,因此不会增加应用查询路径的负载。但复制槽会强制保留未消费的 WAL,可能增加存储开销。此外,Debezium 自动创建发布(publication),用于指定哪些表参与逻辑复制,默认覆盖所有表,也可按需配置。
Q&A
Debezium 是什么?它如何实现 PostgreSQL 的 CDC?
Debezium 是一个分布式 CDC(变更数据捕获)平台,通过追踪 PostgreSQL 的 WAL(预写日志)来捕获行级变更,并将每个变更作为事件发布到 Kafka。它作为 Kafka Connect 的连接器运行,无需自定义代码。
如何用 Docker 搭建 Debezium 的本地环境?需要哪些服务?
需要五个服务:PostgreSQL(设置 wal_level=logical)、Zookeeper、Kafka、Kafka Connect(安装 Debezium 插件),以及可选的 Confluent Control Center。可以使用 docker-compose.yml 来定义这些服务。
注册 Debezium PostgreSQL 连接器的关键配置有哪些?
关键配置包括:connector.class 设为 io.debezium.connector.postgresql.PostgresConnector,database.hostname、database.port、database.user、database.password、database.dbname 指定数据库连接信息,topic.prefix 设置主题前缀,plugin.name 设为 pgoutput,slot.name 指定复制槽名称,publication.autocreate.mode 设为 all_tables,snapshot.mode 设为 initial。
Debezium 事件的结构是怎样的?各个字段代表什么?
每个事件包含 before(变更前的行状态)、after(变更后的行状态)、source(元数据,如数据库、表、事务 ID、WAL 位置)、op(操作类型)和 ts_ms(处理时间戳)。before 在快照和 INSERT 时为 null,DELETE 时 after 为 null。
Debezium 如何保证同一行的变更顺序?
Debezium 使用表的主键作为 Kafka 消息的键,这样同一行的所有事件都会发送到同一个 Kafka 分区,从而保证该行变更的顺序性。
Debezium 对 PostgreSQL 性能有什么影响?
Debezium 读取 WAL 而不是表,因此对应用查询路径没有影响。唯一的开销是 WAL 保留:PostgreSQL 必须保留复制槽尚未消费的 WAL 段,这可能会增加存储和 I/O 开销。