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

PostgreSQL 逻辑复制与 CDC:Publication、Slot、顺序和 Schema

PostgreSQL 的逻辑复制(logical replication)经常被同时用于两类场景:

  1. 数据库到数据库的复制:例如把源库中的部分表复制到报表库、搜索库或另一个 PostgreSQL 实例。
  2. 变更数据捕获(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.orderspublic.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_lsnconfirmed_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 发布 UPDATEDELETE,不等于所有表都能成功复制更新和删除。至少要同时检查:

  1. 表是否有主键或合适的 Replica Identity;
  2. 目标端是否存在可匹配的行;
  3. 目标端是否发生了本地修改;
  4. CDC 事件是否携带了足够的旧值;
  5. 下游是否能处理删除事件。

八、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 列,不能直接假定所有历史数据和复制事件都满足约束。通常需要:

  1. 先增加可空列;
  2. 设置或填充默认值;
  3. 回填历史数据;
  4. 验证没有空值;
  5. 最后再收紧约束。

但即使数据库端演进成功,仍要验证 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;

如果保留量持续增加,需要先判断:

  1. 消费者是否停止;
  2. 网络是否中断;
  3. 消费者是否卡在某个大事务;
  4. 下游写入是否变慢;
  5. Slot 是否已经达到保留上限;
  6. 是否需要暂停源端写入或重新初始化。

不能只执行 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,或者只监控一个消费延迟数字,都不足以证明系统能够可靠地传递数据库变化。


系列导航与关联阅读

官方资料

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