跳至内容

29.1 逻辑复制原语

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

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

29.1.1 publication、subscription 与 replication slot

三个对象,三种职责

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

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

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

一个表可以属于多个 publication,一个 publication 可以有多个 subscriber。它可以限定:

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 TABLESFOR 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:

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_namefailover 属性和源端实际 slot 必须一致。

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

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

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_statusreservedextendedunreservedlost
  • safe_wal_size:在配置有 slot WAL 上限时,距离可能 lost 尚可写多少 WAL;
  • invalidation_reasonwal_removedrows_removedwal_level_insufficientidle_timeout 等失效原因;
  • failover / synced:是否为可同步到 standby 的 failover slot、是否由上游同步而来。

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

生命周期必须成对

常规生命周期:

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

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

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

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

SELECT pg_drop_replication_slot('pg36_shop_slot');

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

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

源端连接角色至少需要:

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 的角色需要数据库 CREATEpg_create_subscription 权限。 默认 run_as_owner = false 时,apply 对每张表切换为目标 table owner;subscription owner 需要能 SET ROLErun_as_owner = true 看似省事,却让目标 table owner 有机会 通过 trigger 等对象以 subscription owner 权限执行代码,除非安全边界极其明确,不应 把它当默认解法。

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

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 初始同步、流式变更与复制身份

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

逻辑复制的主路径是:

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

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

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含义
iinitialize
dcopying data
ftable copy finished
ssynchronized
rready,进入正常复制

查询:

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 的第一次采样看到两张表分别处于 sd,约 0.55 秒后才都进入 r。如果只看 subscription 已存在或 apply worker 有 PID,会过早宣布 initial copy 完成。

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

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 是两套选择

假设:

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 后,还需要:

ALTER SUBSCRIPTION pg36_shop_sub REFRESH PUBLICATION;

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

REPLICA IDENTITY 回答“改哪一行”

INSERT 只需把新值写入目标;UPDATEDELETE 必须携带足以在目标定位旧行的身份。 默认身份是 primary key:

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
);

选择顺序通常是:

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

使用另一个索引:

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 时:

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 拓扑本身就是数据模型的一部分,应进入设计评审。

监控正在发生什么

目标端:

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';

源端:

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、序列、大对象与冲突边界

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

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

source shop.orders
  -> target shop.orders

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

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 应按兼容顺序编排。例如增加一个源端马上会写入的新列:

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 仍为:

last_value = 1
is_called  = false

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

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

示意:

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 是表列,可以复制;两者不要混淆。还应逐项盘点:

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 报错:

insert_exists
update_exists
multiple_unique_conflicts
permission / RLS / other constraint errors

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

第二类可能不停止:

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 为零,仍不 能证明双端相等。

本章正式注入:

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

结果:

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

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

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

native apply errors / conflict statistics
AND
continuous semantic reconciliation

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

trigger 与 apply 语义

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

迁移前应盘点:

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 合同。

进一步阅读:


返回本章目录 · 下一节:CDC 与复制槽治理 · 查看全书目录 · 查看索引中心

最后更新于