# 逻辑复制原语

LLMS 索引： [llms.txt](/llms.txt)

---

逻辑复制不是“把 WAL 发到另一台机器”。物理复制重放块级变化，逻辑复制则在源端把
WAL 解码为表和 tuple 的变化，经 publication 过滤后交给下游 apply。正因为中间已经
进入逻辑层，源和目标可以是不同 major、不同平台，甚至有不完全相同的物理布局；也正
因为只传逻辑 DML，schema、sequence 和许多数据库对象不会自动跟随。

先把原语分清，后面的迁移状态机才不会建立在错误假设上。

## 29.1.1 publication、subscription 与 replication slot {#item-29-1-1}

### 三个对象，三种职责

| 对象 | 所在端 | 持久状态 | 主要职责 |
|---|---|---|---|
| publication | publisher database | table/column/row/operation 集合 | 定义发送哪些逻辑变化 |
| logical slot | publisher cluster + 单一 database | restart/confirmed LSN、xmin/catalog_xmin | 保存一个消费流的保留与确认边界 |
| subscription | subscriber database | conninfo、publication、slot、apply 选项 | 拉取、初始同步并应用到本地表 |

publication 不是消息队列，也不会因为创建就开始发送。它只是某一个 database 中的
变更集合：

```sql
CREATE PUBLICATION pg36_shop_pub
FOR TABLE shop.customers, shop.orders;
```

一个表可以属于多个 publication，一个 publication 可以有多个 subscriber。它可以限定：

```sql
CREATE PUBLICATION paid_orders
FOR TABLE shop.orders (
  order_id, customer_id, status, amount, updated_at
)
WHERE (status = 'paid')
WITH (
  publish = 'insert, update, delete',
  publish_generated_columns = none
);
```

但需要记住几个条件：

- publication 名称只需在当前 database 内唯一；
- `FOR ALL TABLES` 和 `FOR TABLES IN SCHEMA` 会自动纳入未来对象，权限和变更半径更大；
- column list 必须覆盖 replica identity，才能发布 `UPDATE`/`DELETE`；
- `publish` 控制持续 DML，不控制初始数据复制；
- row filter 对 `TRUNCATE` 无效；
- temporary、unlogged、foreign table、view 和 materialized view 不能加入 publication；
- `CREATE PUBLICATION` 不创建 slot，也不产生网络连接。

subscription 位于目标 database：

```sql
CREATE SUBSCRIPTION pg36_shop_sub
CONNECTION
  'host=source.example port=5432 dbname=shop
   user=logical_reader password=REDACTED
   application_name=pg36_shop_sub
   options=-crow_security=off'
PUBLICATION pg36_shop_pub
WITH (
  copy_data  = true,
  create_slot = true,
  enabled    = true,
  slot_name  = 'pg36_shop_slot',
  streaming  = parallel,
  binary     = false,
  run_as_owner = false
);
```

正常情况下，这条命令同时：

1. 在目标 catalog 创建 subscription；
2. 连接源端；
3. 在源端创建持久 logical slot；
4. 启动 apply worker；
5. 为待同步表启动受资源上限约束的 table sync workers。

创建远端 slot 时，`CREATE SUBSCRIPTION` 不能放进显式事务块。若源和目标只是同一个
PostgreSQL cluster 内的不同 database，同一命令一边等待 slot 创建、一边等待自身事务
提交，可能挂住；官方做法是先单独创建 `pgoutput` slot，再用
`create_slot = false` 绑定。跨集群也可以预建 slot，但 subscription 的 `slot_name`、
`failover` 属性和源端实际 slot 必须一致。

### slot 保存的不是一份数据副本

logical slot 是“从哪个位置继续生成一个 database 的变化流”的持久状态：

```sql
SELECT slot_name,
       plugin,
       slot_type,
       database,
       active,
       active_pid,
       xmin,
       catalog_xmin,
       restart_lsn,
       confirmed_flush_lsn,
       wal_status,
       safe_wal_size,
       invalidation_reason,
       failover,
       synced
FROM pg_replication_slots
WHERE slot_name = 'pg36_shop_slot';
```

关键字段：

- `restart_lsn`：仍可能需要的最老 WAL 边界；
- `confirmed_flush_lsn`：logical consumer 已确认接收的位置；
- `xmin` / `catalog_xmin`：仍需保留的普通行版本和系统目录版本；
- `active` / `active_pid`：当前是否有一个消费者；
- `wal_status`：`reserved`、`extended`、`unreserved` 或 `lost`；
- `safe_wal_size`：在配置有 slot WAL 上限时，距离可能 lost 尚可写多少 WAL；
- `invalidation_reason`：`wal_removed`、`rows_removed`、
  `wal_level_insufficient`、`idle_timeout` 等失效原因；
- `failover` / `synced`：是否为可同步到 standby 的 failover slot、是否由上游同步而来。

slot 名称在整个 cluster 中唯一，但 logical slot 只关联一个 database。不同 CDC
consumer 通常需要不同 slot；两个消费者轮流使用同一 slot，不会各自得到完整历史。
同一时刻也只能有一个 receiver 消费它。

### 生命周期必须成对

常规生命周期：

```text
CREATE SUBSCRIPTION
  -> remote slot created
  -> initial sync slots created and dropped
  -> main slot continuously advances
  -> DROP SUBSCRIPTION
  -> remote main slot dropped
```

若目标端必须删除 subscription，而源端暂时不可达，不能假设远端 slot 也消失了。
PostgreSQL 允许先解除关联：

```sql
ALTER SUBSCRIPTION pg36_shop_sub DISABLE;
ALTER SUBSCRIPTION pg36_shop_sub SET (slot_name = NONE);
DROP SUBSCRIPTION pg36_shop_sub;
```

随后必须在源端按名称、database、plugin、active 状态精确检查并删除孤儿 slot：

```sql
SELECT pg_drop_replication_slot('pg36_shop_slot');
```

这不是通用“清理所有 inactive slot”脚本。一个 inactive slot 可能只是计划内停机的
消费者；只有迁移 owner、保留策略和恢复点都确认不再需要时才能删除。

### 最小权限不是一个万能复制账号

源端连接角色至少需要：

```text
LOGIN
REPLICATION
pg_hba.conf / network allow
CONNECT on publisher database
USAGE on published schemas
SELECT on published tables for initial copy
```

若该角色不是 superuser 或 `BYPASSRLS`，publisher RLS policy 可能参与 initial copy
和 row filter。对不信任所有 table owner 的复制域，conninfo 中
`options=-crow_security=off` 会在后来出现 RLS 时停止，而不是悄悄按 policy 过滤。

目标端创建 subscription 的角色需要数据库 `CREATE` 与 `pg_create_subscription` 权限。
默认 `run_as_owner = false` 时，apply 对每张表切换为目标 table owner；subscription
owner 需要能 `SET ROLE`。`run_as_owner = true` 看似省事，却让目标 table owner 有机会
通过 trigger 等对象以 subscription owner 权限执行代码，除非安全边界极其明确，不应
把它当默认解法。

本章实验还遇到一个很实际的 Pigsty 边界：当前 HBA 用 `+dbrole_readonly` 分类业务
连接。临时账号被授予：

```sql
GRANT dbrole_readwrite TO dbuser_pg36source
  WITH INHERIT FALSE, SET FALSE;

GRANT dbrole_readonly TO dbuser_pg36repl
  WITH INHERIT FALSE, SET FALSE;
```

这只让 HBA 的角色成员测试匹配，不让登录角色继承或 `SET ROLE` 到平台角色。把网络
分类与对象权限拆开，既能连接，也不把“能通过 HBA”误写成“自动拥有业务表权限”。

## 29.1.2 初始同步、流式变更与复制身份 {#item-29-1-2}

### 初始快照与主 apply 流怎样汇合

逻辑复制的主路径是：

```text
publisher backend writes WAL
  -> walsender runs logical decoding
      -> pgoutput emits protocol messages
          -> subscriber apply worker maps qualified table/columns
              -> target transaction commits
```

已有数据不能从未来的 WAL 中凭空恢复，所以每张表还经历：

```text
consistent publisher snapshot
  -> table synchronization worker COPY
      -> per-table temporary synchronization slot
          -> replay changes committed during COPY
              -> hand control to main apply worker
```

状态在 `pg_subscription_rel`：

| code | 含义 |
|---|---|
| `i` | initialize |
| `d` | copying data |
| `f` | table copy finished |
| `s` | synchronized |
| `r` | ready，进入正常复制 |

查询：

```sql
SELECT s.subname,
       r.srrelid::regclass AS relation,
       r.srsubstate,
       r.srsublsn
FROM pg_subscription_rel AS r
JOIN pg_subscription AS s ON s.oid = r.srsubid
ORDER BY s.subname, relation;
```

本章正式 run 的第一次采样看到两张表分别处于 `s` 与 `d`，约 0.55 秒后才都进入
`r`。如果只看 subscription 已存在或 apply worker 有 PID，会过早宣布 initial copy
完成。

`r` 也只证明 PostgreSQL 的 table sync 状态。正式验收还比较：

```text
customers rows + ordered digest
orders rows + ordered digest + amount sum
status distribution
orphan orders
negative amounts
invalid statuses
```

结果为 5,000 customers、20,000 orders，双端 logical manifest 相等。

### publication 动作与 initial copy 是两套选择

假设：

```sql
CREATE PUBLICATION insert_only
FOR TABLE shop.orders
WITH (publish = 'insert');
```

`publish = 'insert'` 只限制后续 DML。默认 `copy_data = true` 时，已有行仍会被初始复制。
因此不能用 publication operation list 推断目标基线只含某类事件。

row filter 与 column list 有各自的版本规则；跨版本迁移必须按源、目标最低版本检查。
多个 publication 若以不同 column list 重叠发布同一张表，并不是一个可随意叠加的投影
系统。变更 publication 后，还需要：

```sql
ALTER SUBSCRIPTION pg36_shop_sub REFRESH PUBLICATION;
```

新表才进入 subscription catalog；是否 `copy_data` 要显式决定。`REFRESH` 不是 DDL
迁移，也不会为目标创建缺失表。

### `REPLICA IDENTITY` 回答“改哪一行”

`INSERT` 只需把新值写入目标；`UPDATE` 和 `DELETE` 必须携带足以在目标定位旧行的身份。
默认身份是 primary key：

```sql
SELECT n.nspname,
       c.relname,
       c.relreplident,
       i.indexrelid::regclass AS identity_index
FROM pg_class AS c
JOIN pg_namespace AS n ON n.oid = c.relnamespace
LEFT JOIN pg_index AS i
  ON i.indrelid = c.oid
 AND i.indisreplident
WHERE c.oid IN (
  'shop.customers'::regclass,
  'shop.orders'::regclass
);
```

选择顺序通常是：

```text
stable primary key
  > suitable unique index
      > carefully-reviewed REPLICA IDENTITY FULL
          > no UPDATE/DELETE publication
```

使用另一个索引：

```sql
ALTER TABLE shop.orders
  REPLICA IDENTITY USING INDEX orders_external_id_key;
```

显式 identity index 仍必须满足 unique、immediate、非 partial、列非空等约束。另一条
容易混淆的 PostgreSQL 18 能力是：源端使用 `FULL` 时，目标端可用符合条件的 B-tree
或 hash 候选索引辅助查找；候选不能是 partial，左侧首字段必须是表列而非表达式。
源端 identity 不是 `FULL` 时，目标端也必须有由相同或更少列组成的可用 identity。

没有合适 key 时：

```sql
ALTER TABLE legacy_events REPLICA IDENTITY FULL;
```

这会发送整个旧行。PostgreSQL 18 可以在目标使用满足条件的索引寻找行，但若没有，
每个 update/delete 都可能退化为昂贵查找；某些没有默认 B-tree/hash operator class 的
类型也会限制 apply。`FULL` 是兼容手段，不是免设计主键的奖励。

若 publication 包含 update/delete，而表仍是 `NOTHING` 或默认身份但没有主键，错误会
发生在 publisher 写入路径，而不是等迁移结束才发现。

### transaction order 的保证边界

同一个 subscription 内，subscriber 按 publisher 的提交顺序应用，保持该 stream 的
transactional consistency。这个保证不等于：

- 多个独立 subscription 之间存在全局顺序；
- 外部 CDC sink 的 HTTP、文件或消息副作用与 PostgreSQL commit 原子；
- 目标本地写与源端写自动合并；
- 源数据库之间的事务能组合为一个全局事务；
- 网络恢复后永远不会重发近期消息。

若为了吞吐把相关表拆进多个 subscription，原先同事务的外键或业务原子性也可能被拆开。
publication/subscription 拓扑本身就是数据模型的一部分，应进入设计评审。

### 监控正在发生什么

目标端：

```sql
SELECT subid, subname, worker_type, pid, leader_pid, relid::regclass,
       received_lsn, last_msg_send_time, last_msg_receipt_time,
       latest_end_lsn, latest_end_time
FROM pg_stat_subscription
WHERE subname = 'pg36_shop_sub';

SELECT *
FROM pg_stat_subscription_stats
WHERE subname = 'pg36_shop_sub';
```

源端：

```sql
SELECT application_name,
       state,
       sent_lsn,
       write_lsn,
       flush_lsn,
       replay_lsn,
       reply_time
FROM pg_stat_replication
WHERE application_name = 'pg36_shop_sub';
```

目标显示 apply/sync worker 和接收时间，源端显示 walsender 看到的反馈；两端再用 slot
确认保留边界。不要用 `now() - last_msg_receipt_time` 一项充当“业务复制延迟”：源端
空闲时没有新消息，时间会变大但数据并未落后。更可靠的门禁是产生一个已提交 marker
LSN，等待 `confirmed_flush_lsn` 到达或越过它，并同时验证目标数据。

## 29.1.3 DDL、序列、大对象与冲突边界 {#item-29-1-3}

### DDL 不复制，目标 schema 必须先存在

原生逻辑复制按 fully-qualified table name 匹配目标：

```text
source shop.orders
  -> target shop.orders
```

不会自动创建 schema、table、type、extension、function、constraint、index、owner、
privilege 或 RLS policy。初始 schema 可用受版本控制的 migration 或：

```bash
pg_dump --schema-only --no-owner --no-privileges \
  --dbname="$SOURCE_URL" |
psql --set=ON_ERROR_STOP=1 --dbname="$TARGET_URL"
```

但不能把这条管道当成无需审查的最终方案。`pg_dump` 输出中可能包含 extension、
owner、tablespace、security label、event trigger 依赖；源目标 major 不同时，还需要用
目标版本工具与官方兼容路径评审。

持续 DDL 应按兼容顺序编排。例如增加一个源端马上会写入的新列：

```text
1. target add compatible nullable/defaulted column
2. verify target apply still healthy
3. source add column
4. publication/column list refresh if needed
5. deploy application writes
6. backfill / validate / tighten constraints
```

如果先改源，新的 tuple 已进入 stream，而目标 schema 还无法接收，apply 会报错并停止；
DDL 后来补齐通常能恢复，但期间 slot 继续保留 WAL。

目标 schema 不必字节级相同：

- column 按名称匹配，顺序可不同；
- 目标可有额外列，缺失输入时使用 default；
- 文本模式下，类型只要源文本表示可被目标输入函数接受即可；
- binary 模式要求更严格，跨架构、跨 major 和类型 send/receive 兼容性必须单独验证。

“允许不同”不表示“任意不同都安全”。目标额外 default、trigger、constraint 或
generated expression 可能改变语义或让 apply 失败。

### sequence 不随 identity 值推进

表行中的 identity/serial 数值会被复制，sequence 对象的 `last_value` 不会。本章正式
run 在目标已有 900,000 的最大 `order_id` 时，目标 sequence 仍为：

```text
last_value = 1
is_called  = false
```

如果此刻允许目标接收默认 identity 写入，第一笔就可能生成已经存在的 key。正确顺序是：

```text
source write fence proven
  -> final stream marker acknowledged
      -> source/target manifests equal
          -> synchronize every sequence
              -> verify next value exceeds data high-water mark
                  -> enable target writes
```

示意：

```sql
SELECT setval(
  'shop.orders_order_id_seq',
  greatest(
    (SELECT max(order_id) FROM shop.orders),
    :source_sequence_last_value
  ),
  true
);
```

本章先把目标 sequence 从 1 推进到 900,000，随后目标 canary 得到 900,001。真实系统还
要考虑：

- sequence cache 中已发出但尚未落表的值；
- 多个 sequence 与非标准 ownership；
- shard/tenant 分段号；
- cycling sequence；
- 应用自行生成 ID；
- 回退后源端是否也需要吸收目标 high-water mark。

因此 Pigsty 迁移脚本支持同步 sequence 并可加 offset，但 offset 是冲突缓冲，不是
替代写围栏。

### large object 和非表对象不复制

PostgreSQL large object（`pg_largeobject`/OID API）不在逻辑复制范围。普通表中的
`bytea` 是表列，可以复制；两者不要混淆。还应逐项盘点：

```text
large objects
materialized views and refresh state
views / functions / procedures
extensions and extension versions
FDW server / user mapping
event triggers
roles / memberships / default privileges
RLS policies
database / role settings
tablespaces
collations and ICU/libc versions
scheduled jobs
LISTEN/NOTIFY consumers
external files and object storage references
```

这些对象需要独立迁移和验收，不能因为两张业务表在复制就默认“整个 database 已搬完”。

### subscriber 不是自动只读副本

subscriber 是普通 PostgreSQL database。应用若能在复制目标表本地写入，会产生两类
后果。

第一类会让 apply 报错：

```text
insert_exists
update_exists
multiple_unique_conflicts
permission / RLS / other constraint errors
```

错误型 conflict 会停止复制，必须修复目标数据/权限，或在理解数据损失后显式跳过远端
事务。跳过不是“重试”，而是声明这笔源端事务不再应用到目标。

第二类可能不停止：

```text
update_missing   -> incoming update skipped
delete_missing   -> incoming delete skipped
local target update later overwritten
target-only change remains forever if source never touches that row
```

PostgreSQL 18 在 `pg_stat_subscription_stats` 记录多类 conflict；部分 origin-differs
信息依赖 subscriber 的 `track_commit_timestamp`。但没有报错或 counter 为零，仍不
能证明双端相等。

本章正式注入：

```text
target order_id=900000 note=target-conflict
source order_id=900000 note=source-authority
```

结果：

```text
confl_insert_exists  0 -> 1
apply_error_count    0 -> 1
apply worker         stopped/retried
```

删除精确的目标冲突行后，worker 重放同一源事务并收敛。随后又只改目标
`order_id=1`，这一次没有等待 apply 报错，而是 16 桶摘要中的 bucket 1 不一致。
按源端权威行修复后，mismatch 才归零。

这说明迁移需要两条独立告警线：

```text
native apply errors / conflict statistics
AND
continuous semantic reconciliation
```

只看其中一条，都会漏掉真实故障。

### trigger 与 apply 语义

持续 apply worker 的 `session_replication_role` 为 `replica`，普通 origin trigger/rule
默认不执行；可显式启用 replica/always trigger。初始同步更像 `COPY`，会触发行级和
语句级 INSERT trigger。于是同一张目标表在 initial copy 与 steady state 可能经过不同
trigger 路径。

迁移前应盘点：

```sql
SELECT n.nspname,
       c.relname,
       t.tgname,
       t.tgenabled
FROM pg_trigger AS t
JOIN pg_class AS c ON c.oid = t.tgrelid
JOIN pg_namespace AS n ON n.oid = c.relnamespace
WHERE NOT t.tgisinternal
ORDER BY 1, 2, 3;
```

不要靠 trigger 在目标重新制造本应从源复制的副作用，也不要让 initial copy 重发邮件、
支付请求或 webhook。外部副作用必须有专门的 replay/idempotency 合同。

进一步阅读：

- [PostgreSQL 18：Logical Replication](https://www.postgresql.org/docs/18/logical-replication.html)
- [PostgreSQL 18：Publication](https://www.postgresql.org/docs/18/logical-replication-publication.html)
- [PostgreSQL 18：Subscription](https://www.postgresql.org/docs/18/logical-replication-subscription.html)
- [PostgreSQL 18：Restrictions](https://www.postgresql.org/docs/18/logical-replication-restrictions.html)
- [PostgreSQL 18：Conflicts](https://www.postgresql.org/docs/18/logical-replication-conflicts.html)
- [PostgreSQL 18：Security](https://www.postgresql.org/docs/18/logical-replication-security.html)
- [PostgreSQL 18：`pg_subscription_rel`](https://www.postgresql.org/docs/18/catalog-pg-subscription-rel.html)

---

[返回本章目录](../) · [下一节：CDC 与复制槽治理](../02/) ·
[查看全书目录](/toc/) · [查看索引中心](/indexes/)
