数据库基础体系 · 第 81/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。

数据库 CDC:日志捕获、Debezium、Schema 演进、顺序和重复消费

在数据库与消息系统之间同步数据,常见做法是定期查询:

SELECT * FROM orders WHERE updated_at > ?;

这种方式看似简单,却难以正确处理以下情况:

  • 同一时间有多条记录更新;
  • 事务尚未提交,或者提交后 updated_at 没有可靠变化;
  • 一次更新把字段改回原值;
  • 删除操作没有留下可查询的行;
  • 查询、发送消息、记录游标之间发生崩溃;
  • 下游消费者重复收到同一条变更;
  • 表结构改变后,旧消息和新消息如何共存。

**CDC(Change Data Capture,变更数据捕获)**的核心思路,是捕获数据库已经提交的变化,而不是反复扫描当前表。本文重点讨论基于数据库日志的 CDC,尤其是 MySQL Binlog、PostgreSQL WAL/逻辑解码,以及 Debezium 如何把这些日志转换为消息事件。


一、CDC 到底捕获什么

1. 变化不是“当前值”,而是“状态转换”

对一张表而言,CDC 关注的是某个事务造成的状态变化:

SbeforeTSafterS_{before} \xrightarrow{T} S_{after}

其中:

  • SbeforeS_{before} 是事务提交前的数据库状态;
  • TT 是一个已提交事务;
  • SafterS_{after} 是事务提交后的状态。

对单行数据,常见变化可以表示为:

INSERT: null      -> after
UPDATE: before    -> after
DELETE: before    -> null

因此一个 CDC 事件通常至少需要包含:

  • 操作类型:cud 等;
  • 主键或唯一定位信息;
  • before 和/或 after
  • 源数据库位置;
  • 事务相关信息;
  • 事件对应的 Schema。

例如,Debezium 风格的事件可以抽象为:

{
  "before": {
    "id": 7,
    "status": "NEW",
    "amount": 100
  },
  "after": {
    "id": 7,
    "status": "PAID",
    "amount": 100
  },
  "source": {
    "connector": "mysql",
    "db": "shop",
    "table": "orders",
    "file": "mysql-bin.000123",
    "pos": 456789,
    "row": 2
  },
  "op": "u",
  "ts_ms": 1710000000123
}

这里的 ts_ms 是事件时间或处理相关时间,不能简单等同于数据库提交时间。判断顺序时,应优先使用数据库日志位置、事务边界和连接器提供的源元数据,而不是应用时间戳。

2. CDC 与审计日志不是同一个概念

CDC 通常用于:

  • 同步搜索索引;
  • 构建数仓或湖仓;
  • 缓存失效;
  • 跨系统数据复制;
  • 触发下游业务流程。

审计日志则更强调:

  • 谁执行了什么操作;
  • 请求来源;
  • 业务原因;
  • 变更前后的完整上下文;
  • 长期不可抵赖性。

数据库日志通常能可靠表达“数据发生了什么变化”,但不一定包含“哪个用户因为哪个请求造成了变化”。如果需要后者,应用审计字段、事务上下文或 Outbox 模式仍可能是必要的。


二、为什么捕获数据库日志,而不是轮询表

1. 轮询表存在不可避免的窗口问题

假设消费者按 updated_at 增量读取:

SELECT *
FROM orders
WHERE updated_at > '2024-01-01 10:00:00'
ORDER BY updated_at;

假设两条记录的时间精度只有秒:

id=1, updated_at=10:00:01
id=2, updated_at=10:00:01

如果本轮只处理了 id=1,然后把游标更新为 10:00:01,下一轮使用 > 条件,id=2 就会永久丢失。

改用:

WHERE (updated_at, id) > (?, ?)
ORDER BY updated_at, id

可以修复部分排序问题,但仍无法自然处理:

  • 删除;
  • 时间戳未变化的更新;
  • 事务提交顺序与应用时间顺序不同;
  • 查询结果已读取但消息尚未成功发送时的崩溃;
  • 长事务和并发更新;
  • 数据库主从或读写分离带来的可见性差异。

轮询可以作为补偿或低要求同步方案,但它不是数据库事务日志的等价替代。

2. 日志提供了提交后的变化序列

数据库为了恢复和复制,本来就需要记录变化。CDC 使用这些日志,能够获得:

  1. 事务何时提交;
  2. 事务内包含哪些行变化;
  3. 变化在日志中的位置;
  4. 在某个位置之后继续读取;
  5. 数据库故障恢复后仍可重放的记录。

这里的“日志”不是普通应用日志,而是数据库内部定义的持久化机制:

  • MySQL 主要使用 Binlog 记录复制和恢复所需的逻辑变化;
  • PostgreSQL 先写 WAL,逻辑复制再通过逻辑解码把 WAL 转换为行级变化。

必须区分两个层次:

物理恢复日志:让数据库恢复页面或数据块
逻辑变更日志:让外部系统理解 INSERT/UPDATE/DELETE

PostgreSQL 的 WAL 本身偏物理日志,逻辑解码是把 WAL 还原为逻辑变化的过程。MySQL 的 Row-Based Binlog 则更直接地记录行变化,但它仍然不是面向任意消费者设计的业务事件格式。


三、MySQL CDC:Binlog、事务和行事件

1. 使用 ROW 格式捕获行变化

MySQL Binlog 的常见格式包括:

  • STATEMENT:记录执行的 SQL;
  • ROW:记录受影响的行;
  • MIXED:在两种格式之间由服务器选择。

CDC 通常要求使用 ROW,因为 Statement-Based Logging 存在确定性问题。例如:

UPDATE orders
SET updated_at = NOW()
WHERE status = 'NEW';

如果只记录 SQL,重放时的 NOW()、数据状态和执行环境可能不同。Row-Based Logging 记录的是具体行变化,更适合下游重放。

生产中通常还会明确设置:

[mysqld]
log_bin = mysql-bin
binlog_format = ROW
binlog_row_image = FULL

binlog_row_image=FULL 的意义是:更新事件中包含完整的行镜像。这样 Debezium 或其他消费者更容易得到完整的 before/after 信息。

如果使用 MINIMAL,Binlog 可能只包含:

  • 定位行所需的列;
  • 被修改的列。

这可以降低日志量,但会影响下游构建完整事件、生成反向操作或进行审计的能力。具体事件内容还取决于主键、表结构和连接器配置,不能假定每次 UPDATE 都天然包含完整旧行。

2. MySQL 示例:一条事务对应多个行事件

先创建表:

CREATE TABLE orders (
    id BIGINT PRIMARY KEY,
    status VARCHAR(20) NOT NULL,
    amount DECIMAL(10, 2) NOT NULL
) ENGINE = InnoDB;

执行事务:

START TRANSACTION;

INSERT INTO orders(id, status, amount)
VALUES (1, 'NEW', 100.00);

UPDATE orders
SET status = 'PAID'
WHERE id = 1;

COMMIT;

逻辑上可能看到两条行变化:

INSERT id=1, status=NEW
UPDATE id=1, before.status=NEW, after.status=PAID

但它们属于同一个事务。消费者不能只看到第一条就认为事务已经完成,也不能把事务内的事件当成相互独立的业务事务。

在 Kafka 等消息系统中,Debezium 通常提供事务元数据,使消费者可以识别:

  • 事务标识;
  • 事务事件总数;
  • 当前事件在事务中的序号;
  • 事务开始和结束标记。

是否真正按事务原子性应用,取决于消费者实现。普通消费者逐条处理消息时,仍可能在事务事件之间崩溃。

3. Binlog 位置不是简单的时间戳

MySQL Binlog 的恢复点通常由类似以下信息构成:

binlog file = mysql-bin.000123
position    = 456789

启用 GTID 时,也可以使用 GTID 集合作为恢复位置。位置的本质是日志游标,不是业务时间。

两个事件可能具有:

event A: ts=10:00:01, position=1000
event B: ts=09:59:59, position=1100

不能据此认为 B 先发生。应用时间可能来自客户端,数据库时间可能受时钟、事务执行和提交延迟影响。对同一 Binlog 流,日志位置和事务提交顺序才是更可靠的顺序依据。


四、PostgreSQL CDC:WAL、逻辑解码、Publication 和 Slot

1. WAL 与逻辑解码

PostgreSQL 的 WAL(Write-Ahead Log)首先服务于崩溃恢复和物理复制。WAL 记录足够的信息,使数据库能够在故障后重放已持久化的变化。

逻辑 CDC 需要逻辑解码:

WAL
  -> 逻辑解码
  -> INSERT/UPDATE/DELETE/事务边界
  -> Debezium 或其他消费者

逻辑解码插件负责定义输出格式。PostgreSQL 原生的 pgoutput 用于逻辑复制,也是 Debezium PostgreSQL Connector 常见的插件选择。

2. Publication 决定捕获哪些表

Publication 是 PostgreSQL 中对逻辑复制表范围的声明。例如:

CREATE PUBLICATION app_pub
FOR TABLE public.orders, public.order_items;

也可以按模式或所有表创建,具体能力取决于 PostgreSQL 版本和语句形式。

Publication 解决的是:

哪些表的变化可以被输出

它不解决:

  • 消费者读到哪里;
  • 日志保留多久;
  • 消费失败后如何恢复;
  • 下游是否幂等。

如果新增表没有加入 Publication,CDC 不会自动凭空获得这张表的变化。应明确检查:

SELECT pubname, puballtables
FROM pg_publication;

以及 publication 的表清单。

3. Replication Slot 保存消费位置并阻止 WAL 过早回收

逻辑复制 Slot 记录消费者确认到的日志位置。只要 Slot 还认为某段 WAL 可能被消费,PostgreSQL 就不能随意回收那段 WAL。

查看 Slot:

SELECT
    slot_name,
    plugin,
    slot_type,
    active,
    restart_lsn,
    confirmed_flush_lsn,
    pg_size_pretty(
        pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
    ) AS retained_wal
FROM pg_replication_slots;

主要字段含义:

  • active:当前是否有消费者连接;
  • restart_lsn:仍需保留的较早位置;
  • confirmed_flush_lsn:消费者确认已经处理到的位置;
  • retained_wal:从当前 WAL 位置到保留起点的大致积压量。

这里的风险非常直接:

消费者停止
  -> Slot 不前进
  -> WAL 持续保留
  -> 磁盘空间增长
  -> 最终可能耗尽磁盘

删除 Slot 也不是“清理积压”的无害操作:

SELECT pg_drop_replication_slot('debezium_slot');

删除后,消费者失去原来的恢复位置。重新创建 Slot 后,通常需要重新快照或从其他可用位置重建,不能把它当成普通重置游标。

4. PostgreSQL 的 Replica Identity 决定 UPDATE/DELETE 能否定位旧行

对于 UPDATE 和 DELETE,下游需要知道“哪一行发生了变化”。默认情况下,PostgreSQL 通常使用主键作为 Replica Identity。

查看表属性:

SELECT
    n.nspname AS schema_name,
    c.relname AS table_name,
    c.relreplident
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname = 'orders';

常见值包括:

  • d:default,通常使用主键;
  • n:nothing;
  • f:full,使用整行;
  • i:使用指定的 Replica Identity 索引。

没有主键时,可以指定:

ALTER TABLE orders REPLICA IDENTITY FULL;

FULL 让 PostgreSQL 使用完整旧行帮助定位变化,但会增加日志量,并且在大表、高更新量场景下代价明显。它不是无条件的“最佳设置”。

如果表没有可用标识且 Replica Identity 不足,逻辑复制可能无法为 UPDATE/DELETE 提供足够的旧行定位信息。Debezium 事件中的 before 可能为空、字段不完整,或者复制操作无法按预期工作。CDC 消费者不能假定所有表都天然具有稳定主键。


五、Debezium 的角色:把日志变成可消费事件

1. Debezium 不是数据库日志本身

Debezium 通常运行在 Kafka Connect 生态中:

数据库
  -> Debezium Source Connector
  -> Kafka Connect Worker
  -> Kafka Topic
  -> 下游消费者 / Sink Connector

各层责任不同:

  • 数据库负责生成和保留日志;
  • Debezium 负责读取日志、快照、转换事件;
  • Kafka Connect 负责 Connector 生命周期、配置、offset 管理;
  • Kafka 负责消息持久化、分区和消费位点;
  • 下游消费者负责应用事件、处理失败和幂等。

因此,“Debezium 保证不重复”或“Kafka 保证严格全局顺序”都不是准确表述。每个组件只对自己的语义负责。

2. Debezium 事件的基本结构

Debezium 常见事件包含两层结构:

{
  "before": null,
  "after": {
    "id": 1,
    "status": "NEW"
  },
  "source": {
    "db": "shop",
    "table": "orders",
    "lsn": 123456
  },
  "op": "c",
  "ts_ms": 1710000000000
}

op 通常表示:

op 含义
c create,插入
u update,更新
d delete,删除
r read,快照读取

删除事件的典型形态是:

{
  "before": {
    "id": 1,
    "status": "PAID"
  },
  "after": null,
  "op": "d"
}

如果启用了墓碑消息(tombstone),删除后还可能出现:

{
  "key": 1,
  "value": null
}

这类消息通常与 Kafka Log Compaction 有关:它表达“删除这个 key 的最终保留值”。它不是另一个数据库 DELETE,也不应被消费者误当成带有 before 的行事件。

3. 快照与增量流的交界

新启动的 CDC 连接器通常需要先补齐已有数据,再继续读取实时变化:

快照已有数据
  -> 记录或协调日志位置
  -> 从该位置继续读取增量日志

快照事件常见 op=r,实时 INSERT 则是 op=c。消费者必须接受两者都代表“需要把数据同步到下游”,而不能只处理 c/u/d

快照阶段最容易被误解的是一致性边界。一个正确的快照不能简单地:

SELECT 全表
然后再从“当前时间”读日志

因为两步之间可能发生变化,导致:

  • 重复;
  • 漏数据;
  • 快照看到旧值而增量只看到更新后的值;
  • 增量位置无法与快照读建立可靠关系。

连接器会使用数据库和自身支持的机制协调快照与日志位置,但具体锁、隔离级别、阻塞程度和失败恢复方式受数据库、连接器版本和配置影响。部署前应验证目标版本的快照策略,而不能把“快照完成后读日志”理解为一个没有边界条件的过程。

4. 一个典型的 PostgreSQL Connector 配置骨架

下面是用于说明组件关系的配置片段,具体连接器版本还应根据所采用的 Debezium 发行版校验配置名称:

{
  "name": "shop-postgres",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.port": "5432",
    "database.user": "debezium",
    "database.password": "secret",
    "database.dbname": "shop",
    "topic.prefix": "shop",
    "plugin.name": "pgoutput",
    "slot.name": "debezium_shop",
    "publication.name": "app_pub",
    "table.include.list": "public.orders,public.order_items",
    "snapshot.mode": "initial"
  }
}

前置条件包括:

  1. PostgreSQL 开启逻辑复制所需配置;
  2. 用户具备复制和读取相关权限;
  3. app_pub 已经存在,并包含目标表;
  4. debezium_shop 的生命周期由运维系统管理;
  5. Kafka Connect 能持久化自己的 offset 和 schema history。

配置中的 table.include.list 是连接器筛选,Publication 是 PostgreSQL 服务端筛选。两者都限制范围时,最终捕获集合是它们的交集,而不是并集。


六、Schema 演进:数据格式会随着表结构一起变化

1. Schema 不是“附在 JSON 旁边的注释”

CDC 消息中的 Schema 描述:

  • 字段名;
  • 数据类型;
  • 是否可为空;
  • 精度和小数位;
  • 结构嵌套;
  • 默认值;
  • 字段顺序或兼容性信息。

例如数据库执行:

ALTER TABLE orders
ADD COLUMN paid_at TIMESTAMP NULL;

之后的事件可能含有:

{
  "after": {
    "id": 1,
    "status": "PAID",
    "amount": 100.00,
    "paid_at": "2024-03-10T12:00:00Z"
  }
}

在变更之前产生的旧消息没有 paid_at。因此消费者必须能够同时处理:

旧 Schema + 旧事件
新 Schema + 新事件

2. 增加可空字段通常是兼容演进

新增 nullable 字段时:

ALTER TABLE orders
ADD COLUMN remark VARCHAR(200) NULL;

旧消费者忽略未知字段,通常可以继续运行;新消费者读取旧消息时,把缺少的字段视为 null 或默认值,也通常可行。

但这不是绝对保证,取决于:

  • 序列化格式;
  • Schema Registry 兼容级别;
  • 消费者是否严格拒绝未知字段;
  • 应用反序列化器如何处理缺失字段;
  • 数据库默认值是否被纳入事件。

3. 删除、重命名和收紧约束更危险

以下操作不能简单视为安全:

ALTER TABLE orders DROP COLUMN amount;
ALTER TABLE orders RENAME COLUMN status TO state;
ALTER TABLE orders ALTER COLUMN status SET NOT NULL;

删除字段

旧消费者可能仍依赖该字段。即使新事件不再包含它,历史消息和重放任务仍可能需要它。

重命名字段

数据库的重命名通常在 CDC 层表现为:

旧字段消失 + 新字段出现

连接器不一定能把它识别成“同一字段的别名变化”。消费者可能把它当成删除和新增。

收紧可空性

把 nullable 字段改成 NOT NULL,可能破坏仍在写入旧格式的应用,也可能使回放旧事件时失败。

4. 更稳妥的双字段迁移

status 重命名为 state 为例,不建议直接重命名后要求所有消费者同时升级。可以采用阶段性迁移:

ALTER TABLE orders
ADD COLUMN state VARCHAR(20) NULL;

应用先双写:

status = PAID
state  = PAID

CDC 消费者优先读取 state,没有时回退到 status。待所有消费者迁移完成后,再停止写 status,最终删除旧字段。

这套过程的关键不是“多加一个字段”,而是保证在任意发布窗口内,生产者和消费者的 Schema 交集仍然存在。

5. Debezium Schema History 与消息 Schema Registry 不是一回事

这是经常混淆的两个概念。

Schema History

连接器需要知道数据库表结构如何随时间变化,才能正确解析后续日志。例如 MySQL Binlog 中可能出现:

CREATE TABLE
ALTER TABLE ADD COLUMN
ALTER TABLE MODIFY COLUMN

连接器会维护与自身解析相关的 Schema History。它服务于:

如何解释数据库日志中的字段和类型

Schema History 损坏、丢失或与日志位置不匹配时,连接器可能无法继续解析,而不是简单地产生一个“字段为空”的事件。

Schema Registry

Schema Registry 服务于消息序列化和消费者兼容性检查,例如 Avro、JSON Schema 或 Protobuf。它服务于:

生产者发布什么消息结构
消费者能否读取新旧消息

两者关系是:

数据库结构演进
  -> Debezium 解析所需的内部历史
  -> 消息 Schema 演进
  -> 消费者兼容性策略

配置了 Schema Registry,并不意味着数据库 DDL 可以任意执行;保存了 Debezium Schema History,也不意味着所有消费者都能兼容新消息。

6. MySQL 与 PostgreSQL 的 DDL 捕获边界不同

MySQL Binlog 可以记录 DDL 相关事件,Debezium MySQL Connector 能够据此维护 Schema History,并在配置下发布 Schema Change 事件。

PostgreSQL 的逻辑复制流主要表达表数据变化和事务边界,pgoutput 并不等价于“完整 DDL 审计流”。PostgreSQL Connector 对 DDL 的可见性、字段变更处理和 Schema Change 事件能力不能直接套用 MySQL 的理解。特别是,不能因为 PostgreSQL 表执行了 ALTER TABLE,就假定下游一定收到一条与 MySQL 相同的 DDL 消息。

生产中应把 DDL 迁移工具、数据库版本控制和 CDC 消费者兼容测试结合起来,而不是依赖连接器自动替应用解释所有 DDL。


七、顺序:到底保证哪一种顺序

“CDC 保证顺序”这句话必须补充范围。

1. 顺序至少有四个层次

数据库日志顺序

同一个日志流有自己的位置,例如:

LSN 100 -> LSN 120 -> LSN 150

它表示日志记录顺序。

事务顺序

事务内部通常有多个行变化,事务提交形成边界:

T1: row A, row B, COMMIT
T2: row C, COMMIT

下游若要求事务原子性,需要识别 T1 的完整边界。

Kafka 分区顺序

Kafka 只对同一个 Partition 内的消息提供追加顺序。不同 Partition 之间没有全局顺序:

Partition 0: A1 -> A2
Partition 1: B1 -> B2

只能保证 A1 在 A2 前,不能推出 A2 与 B1 的先后。

业务键顺序

如果订单 order_id=7 的所有变化进入同一个分区:

key=7 -> partition hash(7)

则可以获得该订单键的分区内顺序。若生产者改变分区策略、消息被重新写入其他主题,或者下游合并多个来源,这个保证就不再自动成立。

2. 一个形式化条件

设同一业务键 kk 的事件序列为:

Ek=(e1,e2,,en)E_k = (e_1, e_2, \ldots, e_n)

若需要下游最终状态正确,至少需要:

  1. 所有 eie_i 被投递到同一有序通道;
  2. 通道保持发送顺序;
  3. 消费者按通道顺序应用;
  4. 消费者不会把旧事件覆盖新事件;
  5. 重试不会改变事件的逻辑顺序。

如果这些条件成立,则:

apply(e1);apply(e2);;apply(en)apply(e_1); apply(e_2); \ldots; apply(e_n)

可以复现数据库对该键的状态演进。

但如果事件被分到两个分区:

P0: order=7, status=PAID
P1: order=7, status=SHIPPED

消费者可能先处理 SHIPPED,再处理 PAID,最终状态错误。Kafka 不会依据事件中的 ts_ms 自动重排。

3. 事务内更新同一行的算例

执行:

BEGIN;

UPDATE orders SET status = 'PAID'
WHERE id = 7;

UPDATE orders SET status = 'SHIPPED'
WHERE id = 7;

COMMIT;

CDC 可能产生:

e1: before=NEW,     after=PAID
e2: before=PAID,    after=SHIPPED

若下游按 after 覆盖:

NEW -> PAID -> SHIPPED

结果正确。

若下游先收到 e2,后收到 e1:

NEW -> SHIPPED -> PAID

最终结果错误。

若下游使用源位置判断:

只接受 source_position > last_position

则 e1/e2 的乱序可以被拒绝或延迟;但这要求:

  • 源位置可比较;
  • 位置粒度足够;
  • 每个业务键的事件都在同一可比较序列中;
  • 消费者保存位置与业务写入具有一致性。

因此“增加时间戳排序”通常不是可靠修复。

4. 全局顺序通常代价很高

如果要求所有表、所有业务键的全局顺序,需要单分区或中心化排序。这会降低并行度,并且跨数据库、跨连接器时仍可能没有天然统一的日志序列。

更常见的目标是:

同一表 + 同一主键:有序
同一事务:可识别边界
不同主键:允许并行

这是吞吐量和一致性之间更实际的边界。


八、重复消费:为什么至少一次是常态

1. 典型崩溃窗口

考虑一个 CDC 消费者:

1. 从 Kafka 读取事件 E
2. 写入下游数据库
3. 提交 Kafka offset

如果在第 2 步和第 3 步之间崩溃:

下游已经写入 E
Kafka offset 还没有提交

重启后,Kafka 再次投递 E。

这不是 Kafka 出错,而是因为系统必须在两个独立系统之间协调:

下游数据库事务
Kafka 消费位点

如果它们不共享同一个原子提交协议,就存在二者之间的崩溃窗口。

反过来,如果先提交 Kafka offset,再写下游:

Kafka offset 已提交
下游写入失败

重启后不会重新读取,结果是丢失。

因此常见语义是:

至少一次:可能重复,但尽量不丢

2. Debezium offset 与下游处理不是一个事务

Debezium Source Connector 会维护自己的读取位置,例如 MySQL Binlog 位置或 PostgreSQL LSN。Kafka Connect 还会维护 Source Task 的 offset。

这些 offset 只说明:

连接器认为自己读到哪里

它们不自动说明:

你的业务消费者已经成功更新了搜索索引、缓存或下游数据库

同理,PostgreSQL Slot 的 confirmed_flush_lsn 表示消费者对逻辑复制流的确认进度,也不等于所有业务副作用已经完成。

3. 幂等处理的基本形式

对一个事件 ee,幂等处理要求:

f(f(S,e),e)=f(S,e)f(f(S,e),e)=f(S,e)

其中:

  • SS 是下游状态;
  • ff 是应用事件的操作;
  • ee 是同一事件。

例如,直接执行:

UPDATE orders
SET status = 'PAID'
WHERE id = 7;

同一事件执行两次,结果仍是 PAID,通常具有幂等性。

而下面的操作不是幂等的:

UPDATE account
SET balance = balance + 100
WHERE id = 1;

同一 CDC 事件重复执行两次,会增加 200。

4. 使用事件 ID 去重

可以为事件构造唯一标识,例如:

(source connector, source partition, source offset)

或使用数据库日志位置、事务 ID、行事件序号的组合。必须注意:单独使用主键不够,因为同一个主键会在生命周期内产生很多变化。

下游可以建立处理记录表:

CREATE TABLE processed_cdc_event (
    event_id VARCHAR(300) PRIMARY KEY,
    processed_at TIMESTAMP NOT NULL
);

在同一事务内执行:

BEGIN;

INSERT INTO processed_cdc_event(event_id, processed_at)
VALUES ('mysql-bin.000123:456789:2', CURRENT_TIMESTAMP)
ON CONFLICT (event_id) DO NOTHING;

如果插入影响行数为 0,说明事件已经处理过,可以跳过业务写入:

-- 伪代码逻辑
if inserted_rows == 1:
    apply_business_change()
else:
    skip_duplicate()

必须让“记录事件已处理”和“业务变更”位于同一个下游事务中,否则仍有窗口:

已记录 processed
业务写入尚未提交

5. 版本条件可以阻止旧事件覆盖新事件

如果事件带有可比较的源版本 vv,可以使用条件更新:

UPDATE orders_projection
SET status = ?, source_version = ?
WHERE id = ?
  AND source_version < ?;

这要求:

  • source_version 对同一键单调递增;
  • 所有事件使用同一来源和比较规则;
  • 版本比较不会因文件切换、LSN 格式或重置而失效。

MySQL 的 Binlog file/position、GTID,PostgreSQL 的 LSN,都必须结合具体连接器事件格式处理,不能把不同数据库的版本字段直接混用。


九、删除、墓碑和“当前状态同步”

CDC 事件流通常不是一组可以随意合并的快照,而是状态变化记录。

1. DELETE 必须保留主键

删除后的行已经不存在,因此删除事件至少应保留定位信息:

{
  "before": {
    "id": 7,
    "status": "PAID"
  },
  "after": null,
  "op": "d"
}

如果消费者只接收 after,删除事件会变成“没有内容”,必须从消息 key 获取主键。对没有主键的表,删除同步会更加困难。

2. Kafka tombstone 的用途

假设主题使用 Log Compaction:

key=7, value={...order...}
key=7, value=null

后者告诉 Kafka:

可以删除 key=7 的压缩保留值

业务消费者通常需要处理两种删除相关消息:

  1. op=d:数据库行删除事件;
  2. value=null:Kafka 墓碑消息。

二者语义不同。把墓碑消息当成普通 JSON 反序列化,可能触发空指针或格式错误。


十、故障路径:从数据库到消费者的完整恢复

1. MySQL 日志过期前连接器未恢复

Debezium 停止
  -> Binlog 继续产生
  -> 保留策略删除了旧 Binlog
  -> 连接器需要读取的 position 不存在

连接器无法凭空跳过缺失日志,因为跳过意味着潜在丢失。恢复方案通常是:

  1. 从其他仍保留完整日志的副本恢复;
  2. 重新做快照;
  3. 从可验证的一致性边界重建下游;
  4. 对比源库与下游结果。

不能只修改 offset 到“当前最新位置”然后宣称恢复完成,那会绕过中间变化。

2. PostgreSQL Slot 积压

消费者失联
  -> confirmed_flush_lsn 不前进
  -> WAL 不能回收
  -> 磁盘占用上升

诊断:

SELECT
    slot_name,
    active,
    pg_size_pretty(
        pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
    ) AS retained_wal
FROM pg_replication_slots;

恢复时优先修复消费者连接、权限、网络和解析错误。删除 Slot 是破坏性操作,只有在确认可以重建数据链路并接受重新快照时才使用。

3. 连接器能启动但事件解析失败

这类故障经常与 Schema History 有关:

日志位置仍然存在
连接器也能连数据库
但无法解释某条 Binlog 事件

可能原因包括:

  • Schema History 丢失;
  • Schema History 与 offset 不匹配;
  • 表结构变化未被正确记录;
  • 连接器升级改变了解析行为;
  • 旧日志已超出可恢复的 Schema 上下文。

处理时应保留现场:

  • 连接器日志;
  • 当前 offset;
  • Schema History;
  • 数据库 DDL 记录;
  • 失败事件对应的日志位置。

不要先删除内部主题或直接清空 offset,否则会破坏定位和恢复依据。


十一、一个端到端消费模型

下面用伪代码展示一个支持重复消费和乱序防护的下游事务模型。假设:

  • 每个主键事件被设计为同一 Kafka 分区;
  • event_id 全局唯一;
  • source_version 对同一主键递增;
  • 下游数据库支持事务和唯一约束。
def consume(event):
    event_id = event["event_id"]
    key = event["after"]["id"] if event["after"] else event["before"]["id"]
    version = event["source_version"]
    op = event["op"]

    with database.transaction() as tx:
        inserted = tx.execute("""
            INSERT INTO processed_cdc_event(event_id, processed_at)
            VALUES (%s, CURRENT_TIMESTAMP)
            ON CONFLICT (event_id) DO NOTHING
        """, [event_id])

        if inserted.rowcount == 0:
            return "duplicate"

        if op in ("c", "r", "u"):
            after = event["after"]
            tx.execute("""
                INSERT INTO orders_projection(id, status, source_version)
                VALUES (%s, %s, %s)
                ON CONFLICT (id) DO UPDATE
                SET status = EXCLUDED.status,
                    source_version = EXCLUDED.source_version
                WHERE orders_projection.source_version
                      < EXCLUDED.source_version
            """, [after["id"], after["status"], version])

        elif op == "d":
            tx.execute("""
                DELETE FROM orders_projection
                WHERE id = %s
                  AND source_version < %s
            """, [key, version])

    return "applied"

这个模型有几个重要前提:

  1. event_id 只能用于识别同一事件,不能代替业务键;
  2. source_version 必须具有同一键内的可比较性;
  3. 去重记录与投影写入必须在同一事务中;
  4. DELETE 也要带版本判断,否则旧删除可能覆盖新插入;
  5. r 快照事件不能被误判为无效事件;
  6. 处理失败时事务回滚,offset 不应被确认。

若下游不是支持事务的数据库,例如远程 HTTP API,就不能依赖上述原子性。此时可以使用:

  • 下游幂等接口;
  • 以事件 ID 作为请求幂等键;
  • 持久化待发送任务;
  • 重试队列和人工补偿;
  • 定期对账。

十二、常见误解和对应边界

误解一:CDC 只会产生“最终状态”

错误。CDC 捕获的是变化序列:

NEW -> PAID -> SHIPPED

如果下游只关心当前状态,可以自行做压缩或合并;但审计、计费、库存等场景可能需要每一个中间事件。

误解二:Debezium 消息一定不重复

错误。连接器、Kafka Connect、Kafka 消费者和下游写入之间存在崩溃窗口。除非整个链路有明确、端到端验证过的事务语义,否则消费者应按至少一次投递设计。

误解三:Kafka 分区多,消息就有全局顺序

错误。顺序只在单个 Partition 内成立。要保证同一业务键顺序,必须使用稳定的消息 key 和稳定的分区策略。

误解四:有了 Schema Registry 就能随意改表

错误。Schema Registry 只能约束消息 Schema 的兼容性,不能解决:

  • 数据库 DDL 与连接器解析的兼容性;
  • 快照与增量边界;
  • 旧消费者代码;
  • 删除和重命名字段;
  • 下游数据库迁移。

误解五:PostgreSQL Slot 是一个普通 offset

错误。Slot 还影响 WAL 回收。它既是消费状态,也是数据库磁盘空间风险来源。

误解六:把 offset 调到最新就恢复了

错误。这样做通常意味着放弃中间日志,可能造成静默数据丢失。只有在明确接受重建或丢失窗口,并完成源库与下游对账后,才可以采用这种取舍。


十三、生产诊断应先回答的几个问题

发生延迟、重复、乱序或字段错误时,可以按数据流逐层定位。

数据库层

MySQL:

  • Binlog 是否启用;
  • 是否为 ROW
  • 当前 Binlog 文件、位置或 GTID;
  • 目标日志是否仍然存在;
  • binlog_row_image 是否满足事件需求;
  • DDL 是否改变了连接器理解的表结构。

PostgreSQL:

  • WAL 是否正常产生;
  • Publication 是否包含目标表;
  • Slot 是否存在、是否 active;
  • restart_lsn 与当前 LSN 的差距;
  • 表是否有主键或适当 Replica Identity;
  • 逻辑复制用户权限是否完整。

连接器层

  • 当前 Source offset;
  • 快照是否完成;
  • 是否卡在某个表或某个日志位置;
  • Schema History 是否可读;
  • 连接器是否发生重启或任务重平衡;
  • 错误是连接失败、权限失败、解析失败还是序列化失败。

消息层

  • 主题分区数;
  • 同一个业务键是否始终进入同一分区;
  • Partition offset 是否连续推进;
  • 是否存在重试主题、死信主题或重新发布导致的重复;
  • 墓碑消息是否被消费者正确处理。

消费者层

  • 是否先写下游再提交 offset;
  • 重复事件是否能被唯一约束拦截;
  • 去重记录与业务写入是否同一事务;
  • 是否有版本条件防止旧事件覆盖新事件;
  • 失败重试是否可能改变顺序;
  • 是否存在无法幂等的计数、扣款或库存操作。

十四、如何选择 CDC 的一致性目标

CDC 设计不能只说“实时”和“可靠”,应明确目标:

允许重复吗?
允许丢失吗?
要求同键有序吗?
要求事务原子性吗?
需要捕获 DDL 吗?
需要完整 before 镜像吗?
需要支持历史重放吗?

不同目标对应不同代价:

  • 允许重复但不允许丢失:至少一次 + 幂等消费者;
  • 不允许重复也不允许丢失:需要端到端事务或等价的去重与提交协议;
  • 同键有序:稳定分区键 + 下游顺序处理;
  • 全局有序:减少并行度并承担中心化排序成本;
  • 完整旧值:MySQL 通常倾向 FULL 行镜像,PostgreSQL 可能需要 REPLICA IDENTITY FULL,但两者都增加开销;
  • 可长期重放:必须保留足够的 Binlog/WAL、Kafka 消息和 Schema 历史;
  • DDL 可追踪:需要数据库迁移记录与连接器能力共同保证,不能只依赖行事件。

CDC 最终不是“把数据库改动发到 Kafka”这么简单,而是一条由数据库日志、连接器状态、消息分区、Schema 版本和下游事务共同组成的状态传输链路。

只要明确每一层保存什么状态、在哪个位置确认进度、崩溃后从哪里重放,并为同键顺序和重复消费建立可验证的条件,CDC 才能从一个看似实时的同步工具,变成可恢复、可演进的数据基础设施。


系列导航与关联阅读

官方资料

本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。