数据库基础体系 · 第 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 关注的是某个事务造成的状态变化:
其中:
- 是事务提交前的数据库状态;
- 是一个已提交事务;
- 是事务提交后的状态。
对单行数据,常见变化可以表示为:
INSERT: null -> after
UPDATE: before -> after
DELETE: before -> null
因此一个 CDC 事件通常至少需要包含:
- 操作类型:
c、u、d等; - 主键或唯一定位信息;
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 使用这些日志,能够获得:
- 事务何时提交;
- 事务内包含哪些行变化;
- 变化在日志中的位置;
- 在某个位置之后继续读取;
- 数据库故障恢复后仍可重放的记录。
这里的“日志”不是普通应用日志,而是数据库内部定义的持久化机制:
- 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"
}
}
前置条件包括:
- PostgreSQL 开启逻辑复制所需配置;
- 用户具备复制和读取相关权限;
app_pub已经存在,并包含目标表;debezium_shop的生命周期由运维系统管理;- 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. 一个形式化条件
设同一业务键 的事件序列为:
若需要下游最终状态正确,至少需要:
- 所有 被投递到同一有序通道;
- 通道保持发送顺序;
- 消费者按通道顺序应用;
- 消费者不会把旧事件覆盖新事件;
- 重试不会改变事件的逻辑顺序。
如果这些条件成立,则:
可以复现数据库对该键的状态演进。
但如果事件被分到两个分区:
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. 幂等处理的基本形式
对一个事件 ,幂等处理要求:
其中:
- 是下游状态;
- 是应用事件的操作;
- 是同一事件。
例如,直接执行:
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. 版本条件可以阻止旧事件覆盖新事件
如果事件带有可比较的源版本 ,可以使用条件更新:
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 的压缩保留值
业务消费者通常需要处理两种删除相关消息:
op=d:数据库行删除事件;value=null:Kafka 墓碑消息。
二者语义不同。把墓碑消息当成普通 JSON 反序列化,可能触发空指针或格式错误。
十、故障路径:从数据库到消费者的完整恢复
1. MySQL 日志过期前连接器未恢复
Debezium 停止
-> Binlog 继续产生
-> 保留策略删除了旧 Binlog
-> 连接器需要读取的 position 不存在
连接器无法凭空跳过缺失日志,因为跳过意味着潜在丢失。恢复方案通常是:
- 从其他仍保留完整日志的副本恢复;
- 重新做快照;
- 从可验证的一致性边界重建下游;
- 对比源库与下游结果。
不能只修改 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"
这个模型有几个重要前提:
event_id只能用于识别同一事件,不能代替业务键;source_version必须具有同一键内的可比较性;- 去重记录与投影写入必须在同一事务中;
DELETE也要带版本判断,否则旧删除可能覆盖新插入;r快照事件不能被误判为无效事件;- 处理失败时事务回滚,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 才能从一个看似实时的同步工具,变成可恢复、可演进的数据基础设施。
系列导航与关联阅读
- 系列入口:数据库完整学习路线:从关系模型、事务索引到分布式与向量检索
- 上一篇:时态数据与审计历史:有效时间、系统时间、版本表和可追溯性
- 下一篇:OLTP、OLAP 与 Lakehouse:工作负载、存储布局和数据链路选型
- 延伸:MySQL Redo、Undo 与 Binlog:提交链路、崩溃恢复和一致性
- 延伸:PostgreSQL 逻辑复制与 CDC:Publication、Slot、顺序和 Schema
官方资料
本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论
0 条讨论