数据库基础体系 · 第 100/139 篇。文章以各产品官方稳定版本的公开语义为准;示例会明确引擎、事务与部署边界。
PostgreSQL 逻辑复制与 CDC:Publication、Slot、顺序和 Schema
PostgreSQL 的逻辑复制(logical replication)经常被同时用于两类场景:
- 数据库到数据库的复制:例如把源库中的部分表复制到报表库、搜索库或另一个 PostgreSQL 实例。
- 变更数据捕获(Change Data Capture,CDC):持续读取数据库提交的行级变更,投递到消息队列、数据湖、搜索引擎或下游服务。
这两类场景使用的底层机制相近,但目标不同:
- 数据库复制关心目标库能否正确应用变更;
- CDC 更关心事件格式、幂等、重放、消费位点、Schema 演进和下游顺序。
理解 PostgreSQL 逻辑复制,不能只记住“Publication 发布表,Slot 保存位点”。还需要把以下链路连起来:
事务修改表
↓
WAL 记录物理变更
↓
逻辑解码还原为行级变更
↓
Publication 过滤哪些表、哪些操作、哪些行和列
↓
Logical Replication Slot 保存消费进度并防止相关 WAL 被回收
↓
pgoutput 或其他输出插件编码事件
↓
数据库订阅者或 CDC 消费者读取、确认、应用
一、逻辑复制和物理复制解决的问题不同
PostgreSQL 的 物理流复制传输的是 WAL 中描述的数据页或物理恢复所需的信息。备库通常必须是与主库一致的 PostgreSQL 数据目录副本,复制粒度接近实例或集群。
逻辑复制则从 WAL 中解码出逻辑上的数据库操作,例如:
某事务开始
表 public.orders 的 id=10 被插入
表 public.orders 的 id=10 被更新
某事务提交
逻辑复制不要求目标端拥有相同的数据文件布局,因此可以实现:
- 只复制某些表;
- 源库和目标库使用不同的数据组织方式;
- 复制到另一个 PostgreSQL 数据库;
- 由 CDC 程序把变化转换为消息或其他格式。
但“逻辑”并不意味着不依赖 WAL。逻辑复制仍然依赖源库产生的 WAL,以及源库的逻辑解码能力。
1. WAL、LSN 和逻辑解码
WAL(Write-Ahead Log) 是 PostgreSQL 的预写日志。事务修改数据页前,相关恢复信息先写入 WAL。
LSN(Log Sequence Number) 是 WAL 中的位置。它表示一个单调向前增长的日志位置,而不是业务时间戳,也不是某一行数据的版本号。
逻辑解码会读取 WAL,并将物理层面的记录转换成逻辑变更。逻辑解码必须保留事务边界,因为下游不能把一个尚未提交的事务当成已经成功的数据。
因此,一个事务通常会被表示为:
BEGIN
INSERT ...
UPDATE ...
DELETE ...
COMMIT
对于 CDC 来说,事务边界和提交顺序是核心语义;对于数据库复制来说,它们决定目标库何时可以安全地应用一批变更。
二、Publication:定义“发布哪些变化”
1. Publication 是过滤规则,不是消费位点
Publication(发布)是源库中的逻辑对象,用来声明哪些表和哪些操作对逻辑复制可见。
例如:
CREATE PUBLICATION app_pub
FOR TABLE public.orders, public.order_items
WITH (publish = 'insert, update, delete, truncate');
这条命令表达的是:
- 发布
public.orders和public.order_items; - 发布插入、更新、删除和截断;
- 不创建 Slot;
- 不启动任何网络传输;
- 不保存任何消费者进度。
Publication 只回答:
如果有逻辑复制消费者,它允许看到哪些变化?
而 Slot 回答:
这个消费者已经读到哪里?
这是两个独立的概念。
可以查询 Publication:
SELECT pubname, puballtables, pubinsert, pubupdate,
pubdelete, pubtruncate
FROM pg_publication;
查询 Publication 包含的表:
SELECT p.pubname,
n.nspname AS schema_name,
c.relname AS table_name
FROM pg_publication p
JOIN pg_publication_rel pr ON pr.prpubid = p.oid
JOIN pg_class c ON c.oid = pr.prrelid
JOIN pg_namespace n ON n.oid = c.relnamespace
ORDER BY p.pubname, n.nspname, c.relname;
2. Publication 不发布 DDL
Publication 主要发布表数据变化,不会自动把以下操作作为普通逻辑复制事件同步到目标库:
ALTER TABLE public.orders ADD COLUMN note text;
CREATE INDEX ...;
CREATE TABLE ...;
ALTER TABLE ... ADD CONSTRAINT ...;
因此,源库和目标库的 Schema 管理必须单独完成。后文会详细说明这对 CDC 和数据库订阅分别意味着什么。
Publication 也不是完整的数据库备份。它不会自动复制:
- 表结构定义;
- 权限和角色;
- 索引定义;
- 大多数数据库级配置;
- 序列当前值;
- 未被纳入 Publication 的对象。
3. 按操作类型发布
可以只发布某些操作:
CREATE PUBLICATION order_insert_pub
FOR TABLE public.orders
WITH (publish = 'insert');
这适合“只捕获新增事件”的场景,但必须明确其后果:下游不会知道后续更新和删除。若下游需要构造当前状态,仅发布 insert 通常是不够的。
修改方式:
ALTER PUBLICATION app_pub
SET (publish = 'insert, update, delete');
删除表:
ALTER PUBLICATION app_pub
DROP TABLE public.order_items;
4. 行过滤和列列表
较新的 PostgreSQL 版本支持对表设置行过滤和列列表。例如:
CREATE PUBLICATION active_order_pub
FOR TABLE public.orders
WHERE (status <> 'deleted')
WITH (publish = 'insert, update, delete');
它表达的是发布端过滤条件,而不是目标库上的权限控制。过滤条件在源库侧评估,不能替代应用层的数据脱敏设计。
列列表示例:
CREATE PUBLICATION order_public_fields
FOR TABLE public.orders (id, customer_id, status, updated_at);
这表示只把指定列作为该 Publication 的列变更发布。使用列列表时必须检查:
- 更新和删除是否仍然包含目标端识别行所需的列;
- 下游是否能用不完整的记录重建状态;
- 目标库的表结构和默认值是否匹配;
- 不同 Publication 对同一张表的配置是否造成重复或不一致。
行过滤、列列表和分区表行为具有版本差异,部署时应以目标 PostgreSQL 版本对应文档为准。不能把“创建 Publication 成功”理解为“所有表都一定能被目标端正确应用”。
三、Logical Replication Slot:保存解码状态并约束 WAL 回收
1. Slot 是服务端持久化游标
Replication Slot(复制槽)是 PostgreSQL 服务端保存的一组消费状态。逻辑 Slot 通常与某个数据库绑定,并由一个逻辑解码输出插件读取。
常见创建方式:
SELECT *
FROM pg_create_logical_replication_slot('app_cdc_slot', 'pgoutput');
这里的 pgoutput 是 PostgreSQL 原生逻辑复制输出插件。其他 CDC 工具可能使用自己的输出插件,以生成更适合工具的事件格式。
查询 Slot:
SELECT slot_name,
slot_type,
database,
plugin,
active,
restart_lsn,
confirmed_flush_lsn,
wal_status,
safe_wal_size
FROM pg_replication_slots
WHERE slot_type = 'logical';
不同 PostgreSQL 版本的系统视图列可能有所变化,因此生产脚本不应假设所有版本都具有完全相同的列。
2. restart_lsn 和 confirmed_flush_lsn
两个重要位置的含义不同。
restart_lsn
它表示为了让该 Slot 继续解码,服务器仍可能需要保留的 WAL 起点。只要 Slot 没有推进到足够远,相关 WAL 就不能像普通 WAL 那样被安全回收。
长期不消费的 Slot 可能导致:
pg_wal 持续增长
磁盘耗尽
数据库无法继续写入
confirmed_flush_lsn
它表示消费者向服务器确认已经处理到的逻辑位置。逻辑复制协议中的反馈消息会推动这个位置。
一个简化关系是:
消费者读取事件
↓
消费者将事件应用到目标系统
↓
消费者确认对应 LSN
↓
Slot 的 confirmed_flush_lsn 前进
↓
服务器有机会回收更早的 WAL
如果消费者“刚读到就确认”,但实际还没有写入下游,进程崩溃后可能丢失尚未完成处理的事件。因此确认点必须和消费者自己的持久化状态绑定。
3. Slot 不等于消息队列 offset
Slot 位点和 Kafka offset、数据库表中的处理标记不是同一个东西。
例如:
PostgreSQL Slot 已确认到 LSN 0/500
Kafka 只写入到 offset 100
如果此时消费者先确认了 Slot,再在写 Kafka 前崩溃,重启时 PostgreSQL 可能不会再发送已经确认的事件,而 Kafka 永远缺少这一段数据。
安全顺序通常是:
读取事件
↓
写入下游并持久化
↓
下游写入成功
↓
确认 PostgreSQL Slot
这仍然不自动产生端到端 exactly-once,因为“下游写入成功”和“确认 Slot”之间仍可能发生崩溃。工程上通常依靠:
- 下游幂等键;
- 唯一约束;
- 事务性 offset;
- 可重放日志;
- 去重表;
- 目标系统的原子写入与位点保存。
4. Slot 的生命周期和风险
创建 Slot 后,即使 CDC 程序已经停止,Slot 仍可能保留。删除 Slot:
SELECT pg_drop_replication_slot('app_cdc_slot');
只能删除不再使用的 Slot。删除 Slot 的直接后果是服务器不再为该 Slot 保留 WAL;如果消费者之后重新启动,通常无法从原位置继续,只能重新初始化或从其他备份恢复。
可以通过以下查询发现危险 Slot:
SELECT slot_name,
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
WHERE slot_type = 'logical';
这里的 retained_wal 是当前 WAL 位置和 Slot 保留起点之间的差值,只能作为诊断指标,不能简单等同于“磁盘上一定多占用了这么多空间”。
生产中还必须设置和监控 Slot 的保留上限。过度依赖“定期人工删除 Slot”会把数据库磁盘风险变成事故。
四、从事务到消费者:一条逻辑复制数据流
一个简化的源库到 CDC 消费者流程如下。
步骤一:源库产生事务
BEGIN;
INSERT INTO public.orders(id, customer_id, status)
VALUES (1001, 7, 'created');
UPDATE public.customers
SET last_order_id = 1001
WHERE id = 7;
COMMIT;
WAL 中记录了这些操作以及事务提交信息。
步骤二:逻辑解码识别事务
逻辑解码不会在 INSERT 执行后立即把事件当成已提交事件发送。它需要观察事务提交。
因此,下游通常接收到类似如下的结构:
BEGIN transaction_id=...
INSERT public.orders ...
UPDATE public.customers ...
COMMIT commit_lsn=...
事务中的中间状态可能不会被下游看到。例如:
BEGIN;
INSERT INTO orders(id) VALUES (1);
DELETE FROM orders WHERE id = 1;
COMMIT;
逻辑输出可能包含插入和删除,也可能在某些输出格式中被下游进一步折叠;不能假设“最终表状态相同,所以 CDC 一定没有事件”。CDC 传递的是变更事件,是否合并由输出插件和消费者决定。
步骤三:Publication 过滤
如果 Publication 只包含 orders,则 customers 的更新不会通过该 Publication 发布。
如果 Publication 只发布 INSERT,则事务中的 UPDATE 也不会作为普通更新事件发布。
Publication 过滤不会改变源库事务本身,也不会回滚未发布的表修改;它只改变消费者看到的逻辑输出。
步骤四:输出插件编码
pgoutput 是 PostgreSQL 原生逻辑复制使用的输出格式。它会发送诸如:
- 开始事务;
- Relation 表元数据;
- Insert;
- Update;
- Delete;
- Truncate;
- 提交事务;
- Relation 或 Schema 变化相关的协议元数据。
CDC 工具可能将其转换为 JSON、Avro、Protobuf 或内部事件。例如:
{
"source": {
"db": "app",
"schema": "public",
"table": "orders",
"commit_lsn": "0/5001230",
"txid": 8123
},
"op": "u",
"before": {
"id": 1001,
"status": "created"
},
"after": {
"id": 1001,
"status": "paid"
}
}
这里的 JSON 字段并不是 PostgreSQL 逻辑复制协议的统一标准,而是 CDC 工具自己的事件模型。before 是否存在、LSN 字段如何命名、删除事件如何表示,都必须看具体工具。
五、顺序:哪些顺序有保证,哪些没有
“CDC 保序”是最容易被说得过于绝对的概念。必须区分至少五种顺序。
1. 单个事务内部的顺序
对于同一个事务,逻辑复制保留事务边界,并按逻辑解码得到的顺序发送其变更。
例如:
BEGIN;
UPDATE accounts SET balance = balance - 10 WHERE id = 1;
UPDATE accounts SET balance = balance + 10 WHERE id = 2;
COMMIT;
下游不能把第二个更新应用到第一个更新之前,再声称仍然完全等价。特别是当 CDC 消费者把一个事务拆成多个并发任务时,必须额外维护事务内部顺序,或者等整个事务事件准备好后再提交处理结果。
2. 并发事务之间的提交顺序
设有两个并发事务:
T1: 修改 A ─────────────── COMMIT
T2: 修改 B ── COMMIT
如果 T2 先提交,逻辑复制流通常按提交顺序先发送 T2,再发送 T1:
T2 changes
T2 COMMIT
T1 changes
T1 COMMIT
关键点是:执行开始时间不等于提交顺序。
更形式化地说,设事务集合为:
T = {T1, T2, ..., Tn}
定义 commit(Ti) 为事务提交时对应的逻辑顺序。对同一个逻辑复制流,若:
commit(T1) < commit(T2)
则下游应先看到 T1 的提交边界,再看到 T2 的提交边界。
但这个保证有明确边界:
- 只适用于同一个源 PostgreSQL WAL 流;
- 只适用于同一个 Slot 或由同一复制流导出的顺序;
- 不等于多个 CDC worker 的实际完成时间顺序;
- 不等于多个数据库、多个源集群之间存在全局顺序。
3. LSN 顺序不是业务时间顺序
LSN 大小反映 WAL 位置:
LSN1 < LSN2
通常表示 LSN1 在 WAL 中位于 LSN2 之前,但它不表示:
- 业务事件时间更早;
- 应用请求更早到达;
- 用户看到的时间更早;
- 两行之间存在因果关系。
例如,应用在 10:00 创建订单,但由于事务长时间未提交,另一笔事务在 10:01 提交并可能先被逻辑复制发送。CDC 的排序依据是数据库提交和 WAL 逻辑顺序,而不是应用写入时间。
4. 表内顺序和跨表顺序
如果源库在同一个事务中先写 orders,再写 order_items,下游通常能观察到这个事务内部的顺序和边界。
但如果两个表由两个独立事务写入:
T1: INSERT orders
T2: INSERT order_items
它们之间没有应用层声明的原子关系。即使当前观察到 orders 先到,也不能因此构造业务上的强因果关系。
5. 多分区、多 worker 和消息队列
PostgreSQL 给出的是源复制流的顺序。CDC 系统如果把事件分发到多个分区或多个并发消费者,可能破坏全局到达顺序:
PostgreSQL 顺序:
E1 → E2 → E3
并行处理完成顺序:
E2 → E1 → E3
即使消息系统分区内有序,也只能在同一分区键范围内成立。常见做法是使用业务主键作为分区键,以保证同一实体的事件进入同一分区:
partition_key = order_id
这保证的是同一 order_id 的局部顺序,不是整个数据库的全局顺序。
六、一个完整的顺序算例和反例
假设源库执行以下事务。
T1:
UPDATE inventory SET quantity = quantity - 1 WHERE sku = 'A';
COMMIT;
T2:
UPDATE orders SET status = 'paid' WHERE id = 10;
COMMIT;
如果 T2 先提交,逻辑流可能是:
BEGIN T2
UPDATE orders
COMMIT T2, LSN=0/200
BEGIN T1
UPDATE inventory
COMMIT T1, LSN=0/240
虽然 T1 可能更早开始执行,但下游必须按提交顺序处理。
现在加入两个消费者线程:
线程 1 处理 T2,耗时 500 ms
线程 2 处理 T1,耗时 10 ms
完成顺序变成:
T1 完成
T2 完成
这不代表 PostgreSQL 违反了顺序,而是消费者自己把“读取顺序”和“处理完成顺序”分离了。
再看一个反例:
源库 A 的 LSN:0/500
源库 B 的 LSN:0/600
不能推出:
A 的事件一定早于 B 的事件
因为不同 PostgreSQL 集群的 LSN 命名空间相互独立,两个 LSN 没有可比较的全局含义。
七、Replica Identity:更新和删除如何定位目标行
插入只需要新行数据。更新和删除还需要告诉目标端:
应该修改或删除哪一行?
这由 **Replica Identity(副本标识)**决定。
查看表的副本标识:
SELECT n.nspname,
c.relname,
c.relreplident
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE c.relname IN ('orders', 'customers');
常见取值:
d:default,使用主键;i:使用指定的 Replica Identity 索引;f:full,使用旧行的全部列;n:nothing,不提供用于定位旧行的信息。
最常见的情况是表有主键:
CREATE TABLE public.orders (
id bigint PRIMARY KEY,
status text NOT NULL,
updated_at timestamptz NOT NULL
);
如果没有主键但有合适的唯一索引,可以指定:
ALTER TABLE public.orders
REPLICA IDENTITY USING INDEX orders_business_key_idx;
如果使用:
ALTER TABLE public.orders REPLICA IDENTITY FULL;
更新和删除会携带更多旧行信息,目标端可以通过旧行值匹配记录,但代价可能较高,而且旧行值重复时会产生定位和应用问题。
如果是:
ALTER TABLE public.orders REPLICA IDENTITY NOTHING;
则更新和删除无法提供正常的旧行定位信息。对于需要复制这些操作的场景,这通常会导致逻辑复制或下游应用失败,而不是“自动正确同步”。
一个容易忽略的事实
Publication 发布 UPDATE 和 DELETE,不等于所有表都能成功复制更新和删除。至少要同时检查:
- 表是否有主键或合适的 Replica Identity;
- 目标端是否存在可匹配的行;
- 目标端是否发生了本地修改;
- CDC 事件是否携带了足够的旧值;
- 下游是否能处理删除事件。
八、Schema:逻辑复制不会替你完成结构演进
1. 初始同步和持续复制是两个阶段
原生 PostgreSQL 逻辑订阅通常包含两个阶段:
阶段一:初始表数据同步
阶段二:持续接收初始同步期间及之后的逻辑变更
初始同步不能简单理解为“先执行一次普通 SELECT,再开始复制”。系统必须协调初始快照和变更流,避免漏掉快照期间提交的变化。
在创建订阅时,常见配置类似:
CREATE SUBSCRIPTION app_sub
CONNECTION 'host=source.example dbname=app user=repl password=...'
PUBLICATION app_pub
WITH (
copy_data = true,
create_slot = true,
enabled = true
);
含义是:
- 连接源库;
- 订阅
app_pub; - 默认复制现有数据;
- 默认创建并使用源端 Slot;
- 启用订阅。
这条命令要求目标端的表通常已经存在。原生逻辑复制不会根据 Publication 自动在目标端创建全部表、索引、约束和权限。
2. Schema 不匹配的典型失败
假设源端有:
CREATE TABLE public.orders (
id bigint PRIMARY KEY,
status text NOT NULL,
total numeric NOT NULL
);
目标端只有:
CREATE TABLE public.orders (
id bigint PRIMARY KEY,
status text NOT NULL
);
如果 Publication 发布 total,而目标表没有对应列,应用变更时可能失败。目标端多出一个没有默认值的非空列,也可能导致插入失败。
因此,Schema 变更必须先设计发布端、目标端和 CDC 消费者的兼容关系,而不是等复制报错后再处理。
3. 推荐的加列演进顺序
对于新增可选列,通常采用以下顺序:
第一步:先在目标端增加兼容列
ALTER TABLE public.orders
ADD COLUMN note text;
由于新列允许 NULL,旧事件仍可以应用。
第二步:在源端增加列
ALTER TABLE public.orders
ADD COLUMN note text;
第三步:更新 CDC Schema
如果 CDC 使用 Avro、Protobuf 或 Schema Registry,还要发布兼容的新事件 Schema。
第四步:逐步让应用写入新列
这样,旧消费者即使暂时不知道 note,也不会因为缺列而无法处理原有事件。
如果需要新增 NOT NULL 列,不能直接假定所有历史数据和复制事件都满足约束。通常需要:
- 先增加可空列;
- 设置或填充默认值;
- 回填历史数据;
- 验证没有空值;
- 最后再收紧约束。
但即使数据库端演进成功,仍要验证 CDC 事件的旧格式消费者是否兼容新字段。
4. 删除列、重命名和类型变更
删除或重命名列比加列危险,因为下游可能依赖字段名。
一种安全策略是:
新增新列
双写旧列和新列
让下游切换到新列
停止读取旧列
最后删除旧列
类型变更尤其需要区分三层兼容性:
PostgreSQL 源端类型
逻辑输出插件中的类型
CDC 序列化格式及下游类型
例如,把 integer 变成 bigint,数据库可能允许转换,但下游语言、Schema Registry 或消息消费者未必能无损接受。对于时间、数值精度、枚举和 JSON 字段,也不能仅凭 SQL DDL 判断整个 CDC 链路兼容。
5. DDL 与数据变更的竞态
假设源端执行:
BEGIN;
ALTER TABLE orders ADD COLUMN note text;
INSERT INTO orders(id, status, note) VALUES (1, 'new', 'x');
COMMIT;
目标端如果还没有 note 列,逻辑复制可能在应用插入时失败。因为 DDL 不会自动作为一个可供目标端执行的结构变更同步过去。
Schema 迁移系统必须和 CDC 部署协调。不能只在源端运行迁移脚本,然后等待订阅端“自己跟上”。
九、原生逻辑订阅与自建 CDC 消费者的差异
1. 原生订阅
原生订阅由 PostgreSQL 负责连接源端、读取协议、进行初始同步并在目标端应用变更。
优点:
- 数据库到数据库的语义较完整;
- 事务边界由 PostgreSQL 订阅端处理;
- 不需要自行实现 pgoutput 协议解析;
- 可以直接使用 Publication。
局限:
- 不负责通用消息队列投递;
- 不自动复制 DDL;
- 目标端本地写入可能造成冲突;
- 默认不是双向多主复制方案;
- 目标端的权限、约束、触发器和本地业务逻辑仍可能影响应用。
创建后可以检查订阅状态:
SELECT subname, subenabled, subslotname, subpublications
FROM pg_subscription;
订阅工作进程状态通常可从:
SELECT *
FROM pg_stat_subscription;
观察。不同版本的统计字段可能不同,应结合对应版本文档解释。
2. 自建 CDC 消费者
自建消费者通常通过复制协议连接源库:
复制连接
↓
读取 pgoutput 或其他插件输出
↓
解析 BEGIN / ROW / COMMIT
↓
转换为事件
↓
写入消息系统或下游存储
↓
持久化自己的位点
↓
向 PostgreSQL 发送反馈
它必须处理:
- 网络断开和重连;
- 事务中途断线后的重放;
- Relation 元数据变化;
- 插入、更新、删除和截断;
- 大事务;
- Slot WAL 堵塞;
- 目标写入成功但反馈失败;
- 反馈成功但目标写入未持久化;
- Schema 版本兼容;
- 下游重复消费。
3. 为什么通常只能保证至少一次
假设一个事件的处理过程为:
1. 从 PostgreSQL 读到 E
2. 写入下游
3. 向 PostgreSQL 确认 E 所在位置
如果第 2 步成功、第 3 步失败,消费者重启后可能再次收到 E。
如果顺序改为:
1. 从 PostgreSQL 读到 E
2. 先确认 PostgreSQL
3. 再写入下游
那么第 2 步和第 3 步之间崩溃就可能丢失 E。
因此,跨两个独立系统很难仅凭 Slot 协议得到端到端 exactly-once。可靠设计通常选择“下游先持久化、失败后重放”,也就是接受重复,并让下游幂等。
例如,可以把以下组合设为事件唯一键:
(source_system, database, slot_name, commit_lsn, transaction_sequence)
但具体键是否安全,取决于 CDC 工具如何定义事件序号。仅使用业务主键不够,因为同一行可能发生多次更新;仅使用时间戳也不够,因为时间戳可能重复且不代表复制位点。
十、故障路径:从断线到恢复
场景一:消费者读到事件后立即崩溃
如果消费者尚未发送反馈,Slot 不会前进。重连后,事件会被重新发送。
结果通常是:
不丢失
可能重复
这就是至少一次消费的典型表现。
场景二:消费者已经确认,但下游没有持久化
如果确认已经被源库接受,之后 Slot 可能继续推进并允许 WAL 回收。消费者重启时无法依赖该 Slot 找回已确认的事件。
结果可能是:
丢失
这比重复更难恢复,因此确认必须晚于下游持久化。
场景三:消费者长时间停止
查询:
SELECT slot_name,
active,
restart_lsn,
pg_size_pretty(
pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
) AS retained_wal
FROM pg_replication_slots;
如果保留量持续增加,需要先判断:
- 消费者是否停止;
- 网络是否中断;
- 消费者是否卡在某个大事务;
- 下游写入是否变慢;
- Slot 是否已经达到保留上限;
- 是否需要暂停源端写入或重新初始化。
不能只执行 pg_drop_replication_slot 来“释放空间”。这会丢失该 Slot 尚未交付的逻辑历史。
场景四:目标端应用失败
原生订阅可能因为以下原因停住:
- 目标端缺少表或列;
- 主键或 Replica Identity 不匹配;
- 目标端已有不同数据;
- 约束冲突;
- 类型转换失败;
- 目标端权限不足;
- 触发器或规则改变了应用结果。
正确处理方式是先保留错误现场,确定失败事件和目标状态,再修复 Schema 或数据。直接跳过事务会破坏目标库的一致性;是否允许跳过,取决于业务能否接受目标端形成永久缺口。
十一、配置前置条件和最小示例
源库使用逻辑复制通常需要在配置文件中启用逻辑 WAL,并准备复制连接资源:
wal_level = logical
max_replication_slots = 4
max_wal_senders = 4
然后重启或按参数要求加载配置。还需要在 pg_hba.conf 中允许复制用户通过复制协议连接,并授予所需权限。
例如创建专用用户:
CREATE ROLE cdc_user
WITH LOGIN REPLICATION PASSWORD 'change-this-password';
还需要授予读取 Publication 所涉及对象的权限。权限要求会随 PostgreSQL 版本和使用方式变化,不能只授予 REPLICATION 就假设它自动拥有所有表的读取权限。
创建表和 Publication:
CREATE TABLE public.orders (
id bigint PRIMARY KEY,
customer_id bigint NOT NULL,
status text NOT NULL,
updated_at timestamptz NOT NULL DEFAULT now()
);
CREATE PUBLICATION app_pub
FOR TABLE public.orders
WITH (publish = 'insert, update, delete');
源端写入:
INSERT INTO public.orders(id, customer_id, status)
VALUES (1001, 7, 'created');
UPDATE public.orders
SET status = 'paid',
updated_at = now()
WHERE id = 1001;
DELETE FROM public.orders
WHERE id = 1001;
这个示例的预期是:
- 插入事件被发布;
- 更新事件被发布;
- 删除事件被发布;
TRUNCATE不被发布,因为 Publication 没有启用truncate;- 如果 CDC 消费者重启前没有确认最后一个事件,最后一个事件可能再次出现;
- 如果没有主键,更新和删除还需要额外检查 Replica Identity。
十二、常见误解
误解一:创建 Publication 后 WAL 就会被保存
不准确。
Publication 只是发布规则。真正让服务器保留逻辑解码所需 WAL 的是 Slot。没有消费者或 Slot,Publication 不会提供消费位点。
误解二:Slot 是一个消息队列
不准确。
Slot 是 PostgreSQL 内部的复制状态和 WAL 保留机制,不具备消息队列常见的多消费者、分区、消费组和独立 offset 管理模型。多个消费者通常需要多个 Slot;一个 Slot 被多个独立程序不恰当地共享,可能造成位点互相推进和数据处理混乱。
误解三:LSN 越大,业务事件越新
不完整。
LSN 越大表示日志位置更靠后,但业务时间、请求时间和提交时间不是同一个概念。跨源库的 LSN 更不能比较。
误解四:逻辑复制会同步 DDL
不会自动同步普通 DDL。目标 Schema 必须通过迁移系统、部署脚本或其他控制面单独管理。
误解五:CDC 天然 exactly-once
通常不是。复制流可以在断线后重放,消费者也可能在下游写入和反馈之间崩溃。生产系统应明确选择至少一次、至多一次或经过严格协调的事务性方案,并为重复事件设计幂等处理。
误解六:目标端有相同表名就能复制
不够。还需要检查:
- 表结构和列名;
- 主键和 Replica Identity;
- 类型;
- 默认值;
- 约束;
- 权限;
- 目标端已有数据;
- 触发器和本地写入。
十三、诊断时应把“位点、结构、应用”分开看
逻辑复制故障不要只看一个指标。至少分成三层。
1. Slot 层
SELECT slot_name,
active,
restart_lsn,
confirmed_flush_lsn
FROM pg_replication_slots;
关注:
- Slot 是否存在;
- 是否 active;
confirmed_flush_lsn是否长期不变;restart_lsn与当前 WAL 的距离是否持续增大。
2. Publication 和 Schema 层
SELECT pubname, puballtables, pubinsert, pubupdate,
pubdelete, pubtruncate
FROM pg_publication;
并检查源端和目标端:
\d+ public.orders
关注:
- 表是否被 Publication 包含;
- 操作类型是否启用;
- 目标端是否有对应表和列;
- 主键或 Replica Identity 是否可用;
- 最近是否执行过 DDL。
3. 消费和应用层
关注:
- 最后读取的 LSN;
- 最后确认的 LSN;
- 最后成功写入下游的事件;
- 是否存在重试队列;
- 是否存在目标端约束冲突;
- 是否因为大事务导致事件长时间不可见;
- 消费者重启后是否从正确位置恢复。
尤其要区分:
源端已提交
消费者已读取
消费者已持久化
消费者已确认
目标端已应用
这五个状态不是同一个状态。把它们混为“已经同步”,是排查数据延迟和数据缺失时最常见的根源之一。
十四、生产取舍的核心
PostgreSQL 逻辑复制提供了可靠的数据库级变更流,但它不会替应用系统决定以下问题:
- 哪些数据可以发布;
- 哪些字段需要脱敏;
- 哪些事件必须全局有序;
- 下游如何处理重复;
- Schema 如何兼容演进;
- Slot 堵塞时是否暂停源端写入;
- 目标端冲突是否允许人工修复;
- 重新初始化时如何避免数据缺口。
可以把最终语义概括为:
Publication 决定“看什么”
Slot 决定“从哪里继续看”
逻辑解码决定“如何把 WAL 还原成变更”
事务提交决定“何时对下游可见”
Replica Identity 决定“如何定位更新和删除”
Schema 管理决定“下游能否理解和应用”
消费者确认策略决定“重放、重复和丢失的边界”
因此,设计 CDC 链路时,必须同时验证数据流、事务边界、Slot 生命周期、顺序模型和 Schema 兼容性。只配置一个 Publication,或者只监控一个消费延迟数字,都不足以证明系统能够可靠地传递数据库变化。
系列导航与关联阅读
- 系列入口:数据库完整学习路线:从关系模型、事务索引到分布式与向量检索
- 上一篇:PostgreSQL 全文与模糊检索:tsvector、GIN、Trigram 和相关性
- 下一篇:PostgreSQL 备份恢复:pg_dump、Base Backup、WAL 归档和 PITR
- 延伸:PostgreSQL 复制与高可用:流复制、复制槽、PITR 和故障切换
- 延伸:数据库 CDC:日志捕获、Debezium、Schema 演进、顺序和重复消费
官方资料
本文依据数据库官方文档重新梳理;正文、示例与生产检查清单由 WR BLOG 编写。

评论
0 条讨论