跳转到主要内容

29 移花接木:逻辑复制、迁移与异构同步

逻辑复制能让数据流动,但“数据正在流动”离“迁移已经成功”还很远。

publication 不复制 DDL,sequence 不随表中 identity 值推进,大对象不在复制范围内; subscription 停止后,slot 仍可能继续保留 WAL 与 catalog rows;本地写入 subscriber 既可能造成会报错的 apply conflict,也可能留下完全不报错的静默漂移;在切换窗口里, 源端是否真的停止写、目标序列是否安全、目标新增写如何反向带回,决定了回退是不是一句 空话。

本章把迁移建模为一个有证据、有门禁、可失败关闭的状态机:

inventory and semantic preflight
  -> target schema
      -> initial snapshot
          -> incremental stream
              -> catch-up fence
                  -> source write fence
                      -> sequence / external state sync
                          -> shadow verification
                              -> route switch
                                  -> observation
                                      -> forward or rollback decision

其中每一条箭头都有前置条件、观测、超时、停手方式与恢复路径。不能用最终行数相同跳过 中间状态,也不能把“DNS 已改”当成目标端数据、权限和业务语义已经可用。

本章目标

理解逻辑复制、CDC 和数据搬迁的状态机,用校验与可回退切换完成迁移,而不是把“数据能流动”误当成迁移成功。

读完本章,读者应该能够:

  1. 从 WAL、output plugin、publication、slot、subscription 和 apply worker 解释原生 逻辑复制的数据路径;
  2. 区分 initial table synchronization、持续 DML 与独立的 schema/sequence 迁移;
  3. 为表选择 primary key、REPLICA IDENTITY USING INDEX 或受约束的 FULL
  4. pg_replication_slotspg_stat_subscriptionpg_stat_subscription_statspg_subscription_rel 定位停滞和冲突;
  5. 解释 slot 为什么提供可恢复位点,却不能天然提供 exactly-once 的外部副作用;
  6. 为 CDC sink 设计 event identity、幂等提交、位点原子性与 replay 策略;
  7. 设计 COPY/dump/restore 的并行装载顺序,并保留 rejected rows;
  8. 用行数、摘要、分桶、不变量和抽样分别验证结构与数据;
  9. 把在线迁移拆成 preflight、全量、增量、追平、写围栏、切流、观察和退出;
  10. 区分 rollback、forward repair 与已经跨过不可逆点的前滚;
  11. 为异构系统显式记录类型、精度、时区、排序、约束、顺序和 delete 语义损失;
  12. 用 Pigsty 生成迁移上下文、观察两端集群,但不把生成脚本误认为自动切流授权;
  13. 建立源、目标、验证和管理端点的独立凭据与网络边界;
  14. 在源端主库切换时评审 logical slot failover,而不是默认 subscription 会自动跟随;
  15. 完成一份含原始证据、公开摘要、回退记录和双端精确清理的迁移证据包。

一张图看清复制与迁移

publisher database
  table DML
    -> WAL
      -> logical decoding
        -> pgoutput
          -> publication filter
            -> logical slot
              -> streaming protocol
                -> subscriber apply worker
                  -> same qualified table name

separate migration tracks
  DDL / ownership / privileges
  sequences
  large objects
  extensions and settings
  application routes and credentials
  validation and rollback state

publication 是某一个数据库中的变更集合;subscription 定义下游连接、publication 集合和 apply 行为;一个活动 subscription 通常对应源端一个持久 logical slot,初始 复制还会短暂创建 table synchronization slots。slot 的名称在整个 PostgreSQL cluster 中唯一,但 logical slot 只属于一个 database。

初始复制完成后,pg_subscription_rel.srsubstate = 'r' 只说明表同步状态 ready; 它不证明 sequence、DDL 或业务不变量一致。迁移验收必须把 PostgreSQL 内建状态与 独立 reconciliation 同时纳入。

本章正式实验

本章在两个具有不同 PostgreSQL system identifier 的 Pigsty 沙箱集群之间执行:

source       pg-test / pg36_shop_src
target       pg-meta / pg36_shop_dst
PostgreSQL   18.6 on both sides
wal_level    logical
data         synthetic only
route        private evidence simulation only
production   not touched

正式 run:

run id       7d95ca65-12a7-46c2-8e6c-ad8cbb5336c5
preflight    568c4034-5b7c-4be1-a5e0-48fa189bb782
source sysid 7668025967696967004
target sysid 7668025945980641675

初始复制:

对象 行数 table sync state logical manifest
shop.customers 5,000 r 相等
shop.orders 20,000 r 相等

随后增量执行:

INSERT 500
UPDATE 200
DELETE 100
source marker acknowledged
logical manifests equal

暂停 subscription 后写入 3,000 行:

证据 暂停前 暂停后
slot active false false
confirmed_flush_lsn 不变 不变
retained bytes 227,008 2,867,128

恢复 subscription 后,marker 被确认且双端重新相等。

冲突与静默漂移:

order_id 900000
  target local row + source incoming row
  -> confl_insert_exists: 0 -> 1
  -> apply_error_count:    0 -> 1
  -> remove exact target conflict
  -> replay and converge

order_id 1
  target-only value change
  -> no apply error required
  -> bucket 1 digest mismatch
  -> source-authoritative repair
  -> zero mismatched buckets

切换与回退:

source runtime write attempt     SQLSTATE 42501
target sequence before           1
target sequence synchronized     900000
first target canary              900001
private route history            source -> target -> source
real platform route changed      false
target-only rows reconciled      1
source retained through rollback true
final customers                  5000
final orders                     23402
final manifests equal            true
orphan / negative / invalid      0 / 0 / 0

最后普通删除两端数据库、五个角色、subscription 与 slot;未使用 force drop,未终止 无关会话。29 个预声明反例和 19 个现场证据 mutant 全部被拒绝。公开结果见 migration-run.json

本章目录

29.1 逻辑复制原语

先建立 PostgreSQL 原生模型:哪些对象定义变更集,哪些对象保存位点,初始快照如何追上 主 apply 流,以及 UPDATE/DELETE 怎样定位目标行。最后明确没有进入这条流的内容。

29.2 CDC 与复制槽治理

从内建 subscription 扩展到通用 CDC:区分产生事件、传输、处理、提交外部副作用与推进 位点,解释重放和重复为什么不可避免,并把 inactive slot 转成可预算的 source risk。

29.3 批量装载与数据校验

迁移常常先以批量方式搬运基线。本节讨论 server-side COPY、client-side \copy、 并行与约束顺序,重点是如何把坏行、批次、源文件散列和最终 reconciliation 留下来。

29.4 在线迁移状态机

把工具动作提升为迁移项目:每一 phase 有 entry condition、exit evidence、owner、timeout 与 abort path。读路径和写路径分别验证,切换不是单一时刻,而是逐步关闭不确定性的过程。

29.5 异构同步的语义损失

跨引擎 CDC 不能只看连接器“green”。本节用可声明的 semantic contract 记录类型映射、 精度、排序、事务边界、tombstone、约束和重新处理行为,再用代表性查询验证目标用途。

29.6 多集群迁移环境

用 Pigsty 承担集群、端点、监控和迁移上下文的参考实现。明确 pgsql-migration.yml 生成的是操作手册与脚本,真实写围栏、路由和退出仍需按应用接入 方式设计、审批与执行。

29.7 实战:迁移 pg36_shop

用本章 runner 重演完整双集群实验,阅读私有证据与公开白名单,练习如何从一个失败阶段 安全恢复,而不是只观察成功路径。

推荐学习路线

应用开发者:

29.1 -> 29.3 -> 29.4 -> 29.5 -> 29.7

重点是 schema/sequence 边界、数据校验、应用兼容、读写切换与异构语义。

平台工程师:

29.1 -> 29.2 -> 29.4 -> 29.6 -> 29.7

重点是 slot/WAL、权限、进度、故障恢复、源端主库切换和迁移证据。

两条路线最终必须合流:平台可以保证 stream 可用,无法替应用决定订单状态是否等价; 应用可以定义不变量,无法独自保证 slot、磁盘、网络和切换窗口。

版本与证据权威

本章 PostgreSQL 语义以 18 为基线,正式 run 使用 18.6。逻辑复制能力跨版本变化明显: row filter、column list、binary、streaming、two-phase、conflict statistics、failover slot、generated columns 与 subscription 权限都必须以源、目标实际 major/minor 的官方文档为准。

Pigsty 示例以 4.4 为参考。当前 pgsql-migration.yml 生成迁移上下文、操作手册和脚本, 不会替操作者直接完成真实路由。生成物需要进入变更评审,secret 也不能因为出现在模板里 就进入仓库或公开证据。

核心资料:

实验文件

static/labs/ch29/
├── requirements.json
├── migration-contract.json
├── negative-cases.json
├── topology.mmd
├── lab-contract.md
├── capture.py
├── remote_experiment.py
├── exercise.py
├── validate.py
├── review.py
├── task.sh
└── migration-run.json

执行:

static/labs/ch29/task.sh lint

export PG36_EVIDENCE_DIR=/absolute/private/new-empty/ch29-run
static/labs/ch29/task.sh capture
static/labs/ch29/task.sh exercise
static/labs/ch29/task.sh verify
static/labs/ch29/task.sh review

all 会按相同顺序执行。exercise 会在两个沙箱集群上真实创建 subscription/slot、 生成 WAL 并注入冲突,只能在明确的一次性开发/测试环境运行。

本章验收

只有当迁移证据包能回答以下问题,才算完成:

exact source and target identity
schema / extension / collation / setting diff
publication table and replica identity inventory
initial table synchronization states
source marker LSN and subscriber acknowledgement
slot restart / confirmed LSN, wal_status and safe_wal_size
row count + digest + bucket + business invariant
DDL / sequence / large object / external object plan
apply conflict statistics and repair record
source write-fence proof
route change, owner and rollback point
destination-only write reconciliation
source retention and exit criteria
subscription / slot / credential cleanup
production approval state

“两端 count(*) 相等”不是迁移验收。


上一章:除旧布新:VACUUM、冻结与膨胀治理 · 返回下卷导读 · 下一章:推陈出新:版本升级与回滚策略 · 查看全书目录 · 查看索引中心

29.1 逻辑复制原语

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

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

29.1.1 publication、subscription 与 replication slot

三个对象,三种职责

对象 所在端 持久状态 主要职责
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 中的 变更集合:

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 含义
i initialize
d copying data
f table copy finished
s synchronized
r ready,进入正常复制

查询:

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 与复制槽治理 · 查看全书目录 · 查看索引中心

29.2 CDC 与复制槽治理

内建 logical replication 的目标是 PostgreSQL table,CDC 的目标可能是 Kafka、 对象存储、搜索引擎、缓存、数据仓库或另一个业务服务。两者都从 logical decoding 出发,却在“谁保存位点、何时确认、如何处理重复、怎样表达 schema”上走向不同系统。

本节不绑定某一个 connector,而是建立任何 CDC 实现都绕不开的合同。

29.2.1 逻辑解码、插件与消费位点

从物理 WAL 到业务事件

PostgreSQL WAL 记录的是恢复数据库所需的物理/逻辑底层信息,并不是一个稳定的 JSON 业务接口。logical decoding 负责把一个 database 的持久表变化重组为按事务提交边界 可解释的 change stream:

WAL records
  -> decode relation and tuple changes
      -> reorder by transaction
          -> output plugin
              -> logical replication protocol or SQL decoding API
                  -> consumer

output plugin 决定编码,不决定下游业务语义:

plugin / path 典型用途 重要边界
pgoutput 内建 publication/subscription 二进制 replication protocol,不是给人直接读的文本
test_decoding 测试和理解 logical decoding 测试插件,不是生产业务协议
第三方 JSON/connector plugin 通用 CDC 版本、扩展、schema envelope、升级兼容需独立负责

不要把“输出是 JSON”误认为 schema compatibility 已解决。事件仍需定义:

source identity
database / schema / table
transaction and commit position
operation
replica identity before image
new column values
schema version
event id
producer version

缺少 before image 时,DELETE 怎样定位?TOAST 列未变化时,connector 如何表示?列 rename 是新字段还是同一字段?这些都不是 JSON 格式本身能回答的。

SQL 解码 API 与 streaming protocol

教学或诊断可用:

SELECT *
FROM pg_create_logical_replication_slot(
  'demo_slot',
  'test_decoding'
);

SELECT *
FROM pg_logical_slot_peek_changes(
  'demo_slot',
  NULL,
  100
);

SELECT *
FROM pg_logical_slot_get_changes(
  'demo_slot',
  NULL,
  100
);

peek 不消费,适合观察;get 会推进消费位置。对 pgoutput 应使用 binary API 或 replication protocol,而不是强行调用文本函数。

生产 connector 通常通过 replication protocol:

START_REPLICATION SLOT ...
  -> server sends changes
      -> client sends standby status updates
          -> slot confirmed_flush_lsn advances

“客户端已经读到”“写进本地 buffer”“写到外部系统”“外部事务提交”“向 PostgreSQL 确认”是五个不同瞬间。connector 必须声明哪个瞬间触发 acknowledgement。

snapshot 与 stream 必须无缝衔接

新建 logical slot 时可导出一个 snapshot。这个 snapshot 表示:

snapshot sees database state at boundary B
slot emits all committed changes after boundary B

于是通用 bootstrap 可以:

  1. 创建 slot 并取得 exported snapshot;
  2. 在另一个事务 SET TRANSACTION SNAPSHOT
  3. 全量导出 snapshot 中的基线;
  4. 从 slot 持续消费 B 之后的变化;
  5. 在下游把全量与增量合并。

内建 subscription 将这套协调封装在 table sync workers 中。自行开发 CDC 时,如果先 普通 SELECT 全表、很久以后才创建 slot,会在两者间漏变化;若先创建 slot却不消费也 不预算磁盘,会让 WAL 无界增长。bootstrap 顺序必须成为协议,不是运气。

四个位置不要混为一个“offset”

一个成熟 CDC 管道至少有:

restart_lsn
  publisher still needs WAL from here

confirmed_flush_lsn
  source slot believes consumer confirmed through here

consumer durable checkpoint
  connector can restart from here

sink commit / business watermark
  downstream effects are durable through here

它们应满足与具体协议一致的单调关系。最危险的情况是 connector 先确认 source,再把 buffer 异步写 sink:

confirmed_flush_lsn advances
  -> connector crashes before sink commit
      -> source is allowed to recycle old changes
          -> permanent downstream gap

反过来,sink 已提交而 checkpoint 尚未推进会造成 replay;这通常可以用幂等处理,永久 缺口却无法凭空修复。因此多数 CDC 更愿意接受 at-least-once,而不是冒险提前确认。

plugin、slot 与版本都要进 inventory

建议为每个 consumer 记录:

consumer_id: search-orders-v3
source_cluster: pg-prod
database: shop
slot: cdc_search_orders_v3
plugin: pgoutput
publication: cdc_search_orders
owner: search-platform
schema_contract: order-event-v7
checkpoint_store: kafka-connect-offsets
ack_after: sink_transaction_commit
max_replay_window: 15m
max_slot_retained_wal: 80GiB
rebootstrap_method: snapshot-plus-stream
retirement_ticket: null

slot 本身不知道 consumer owner、SLO 或 sink 状态。没有外部 inventory,inactive slot 只能告诉你“现在没人连”,不能告诉你“可以删”。

29.2.2 至少一次、重复事件与下游幂等

PostgreSQL 已经明确允许 replay

logical slot 是 crash-safe 的,但它的当前位置只在 checkpoint 时持久化。服务器崩溃后, slot 可能回到较早 LSN,于是最近变化再次发送。网络断开、consumer 在 sink commit 后但 ack 前崩溃,也会产生相同结果。

因此正确假设是:

event may be delivered more than once
transaction order is meaningful within a stream
acknowledged history can no longer be requested from that slot

“测试十次没重复”不能升级为 exactly-once 保证。

exactly-once 是端到端属性

若 sink 也是 PostgreSQL,可以在同一个目标事务里同时写业务状态与 dedup ledger:

BEGIN;

INSERT INTO cdc_applied(event_id, source_lsn, applied_at)
VALUES (:event_id, :source_lsn, clock_timestamp())
ON CONFLICT (event_id) DO NOTHING;

-- 只有上一条确实插入时,才应用业务变化。
UPDATE search_projection
SET ...
WHERE ...
  AND :new_event_was_inserted;

COMMIT;

随后才向 source 确认。若事务已提交但 ack 丢失,replay 会命中唯一键,不重复副作用。

但如果副作用跨多个系统:

write database
send email
charge payment
publish Kafka
advance offset

没有一个 PostgreSQL transaction 能原子覆盖全部。应使用:

  • transactional outbox;
  • sink 自身的幂等 key / conditional write;
  • 去重 ledger;
  • 可重建投影;
  • saga/补偿;
  • 对不可重复副作用的业务级 request identity。

不要用“处理成功后更新 offset”一句话掩盖多个 commit 之间的崩溃窗口。

事件身份怎么设计

只用 txid 不够:32-bit XID 会回卷,跨 cluster/database 也会重复。只用表主键也不够: 同一行可以变化很多次。一个实用 envelope 可包含:

source system identifier
database identity
slot / publication contract
commit LSN
transaction identity
change ordinal within transaction
schema contract version

具体 connector 能提供哪些字段取决于协议。核心要求是:

same logical change -> same event id on replay
different logical changes -> different event ids

不要把消费时间戳或随机 UUID 当 event id;replay 时它们会变。

upsert 不自动等于幂等

下面的 sink 写法:

INSERT ... ON CONFLICT (id) DO UPDATE ...

只有在“最后写入覆盖即可”且顺序严格时才近似幂等。它处理不了:

  • 较旧事件晚到,覆盖较新状态;
  • amount = amount + delta 被重复执行;
  • DELETE/tombstone 后旧 UPDATE 复活;
  • 多个 source 同写一个 key;
  • 事件 schema 变化导致部分字段保留旧值;
  • 外部副作用已经发生。

更安全的投影常带 source version:

INSERT INTO order_projection (
  order_id, status, source_commit_lsn, payload
)
VALUES (...)
ON CONFLICT (order_id) DO UPDATE
SET status = excluded.status,
    source_commit_lsn = excluded.source_commit_lsn,
    payload = excluded.payload
WHERE order_projection.source_commit_lsn
    < excluded.source_commit_lsn;

LSN 只在同一 source timeline/合同内有序;跨 source 合并仍需业务 version/vector 或冲突 规则。

transaction boundary 不能随意打散

源事务:

debit account A
credit account B
append ledger

若 connector 把三条 row event 分别确认并让下游实时可见,中途失败可能暴露不平衡状态。 CDC envelope 应保留 begin/commit 或 transaction grouping;sink 要么原子应用整个事务, 要么明确只提供最终一致读模型并隐藏未完成 batch。

大事务还会带来另一组选择:

  • streaming = off:源端完整解码后发送,内存/延迟风险;
  • streaming = on:未提交变化先写 subscriber 临时文件,commit 后应用;
  • streaming = parallel:有 worker 时直接并行 apply,否则回退临时文件。

吞吐设置不能改变“只有源端 commit 后才把结果视为提交”的语义。

poison event 需要隔离,不是静默跳过

遇到无法解析或不满足 sink constraint 的事件:

stop entire stream forever
skip and lose data silently
retry at full speed forever

都不是完整策略。应记录:

event_id: ...
source_position: ...
schema_version: ...
error_class: conversion | constraint | permission | code
first_seen: ...
attempts: ...
raw_payload_ref: encrypted-private-location
owner: ...
decision: repair-and-replay | compensate | approved-skip

dead-letter queue 只是隔离区,不是数据正确性的垃圾桶。任何 approved skip 都要进入 reconciliation,并记录业务影响。

29.2.3 槽停滞、WAL 保留与磁盘风险

inactive 不等于无害

slot 与连接生命周期独立。consumer 下线后:

active = false
confirmed_flush_lsn stops
source continues writing WAL
restart_lsn remains old
pg_wal retained bytes grow
catalog_xmin may also hold catalog tuples

基础查询:

SELECT slot_name,
       database,
       active,
       active_pid,
       inactive_since,
       xmin,
       catalog_xmin,
       restart_lsn,
       confirmed_flush_lsn,
       pg_size_pretty(
         pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)
       ) AS retained,
       wal_status,
       pg_size_pretty(safe_wal_size) AS safe_wal,
       invalidation_reason
FROM pg_replication_slots
WHERE slot_type = 'logical'
ORDER BY restart_lsn NULLS FIRST;

pg_wal_lsn_diff(current, restart_lsn) 是按当前时刻估算该 slot 的 WAL 保留距离,不等于 磁盘上所有 WAL 文件恰好这么大;checkpoint、archive、其他 slot 与 segment 粒度都会 影响实际 pg_wal

把风险换成时间

若近期 WAL 产生率为 rr bytes/s,slot 当前保留 LL,可用于增长的安全空间为 FF,则最粗略的时间预算:

TdiskFr T_{\text{disk}} \approx \frac{F}{r}

若配置 max_slot_wal_keep_size = M,距离 slot 可能在 checkpoint 后失去所需 WAL 的 预算:

TslotMLr T_{\text{slot}} \approx \frac{M - L}{r}

实际告警应使用变化率和低水位:

retained bytes
safe_wal_size
wal_status / invalidation_reason
pg_wal filesystem free
archive health
consumer lag and error rate
catalog_xmin / XID age

只告警 active = false 会在计划维护时产生噪声;只告警磁盘使用率会在 WAL 已快速增长时 太晚。

上限保护的是 source,不保证 consumer 可恢复

max_slot_wal_keep_size 为非负值时,checkpoint 可以允许回收超出上限的 WAL,slot 可能进入 unreserved 乃至 lost。这避免一个遗忘 consumer 无限填满源盘,但代价是 consumer 必须 rebootstrap。

PostgreSQL 18 还提供 idle_replication_slot_timeout:slot inactive 超过时限后可在 checkpoint 被 invalidated。它同样不是“暂停后自动保存到对象存储”。启用前必须明确:

  • 哪类 slot 允许因空闲失效;
  • 谁收到即将失效告警;
  • snapshot/rebootstrap 需要多久;
  • 同步到 standby 的 slot 等豁免/特殊语义;
  • checkpoint 周期带来的执行延迟。

不要为了保护磁盘把上限设得小于正常故障恢复窗口,然后把频繁 lost 当 consumer 问题。

本章停滞实验

正式实验先让 20,000 初始订单和 500/200/100 增量收敛,再执行:

ALTER SUBSCRIPTION pg36_shop_sub DISABLE;

确认目标 subenabled = false、worker 为零、源端 slot active = false 后,在源端插入 3,000 行有界 payload。

结果:

字段 写入前 写入后
active false false
confirmed_flush_lsn 1/9FDBBF30 1/9FDBBF30
retained bytes 227,008 2,867,128
wal_status reserved reserved

这里能确认因 consumer 停滞,确认位点没有推进且保留距离增长。不能用这 2.87 MB 推导 生产每 3,000 行的固定 WAL 成本:其他 database、full-page write、索引、payload、 checkpoint 和并发都会改变 WAL。

恢复:

ALTER SUBSCRIPTION pg36_shop_sub ENABLE;

验收同时要求:

confirmed_flush_lsn >= source marker LSN
both relation states = r
source logical manifest = target logical manifest

“slot active 又变 true”仍不够。

停滞处置顺序

1. identify exact slot and owner
2. confirm source database / plugin / consumer contract
3. capture restart, confirmed, xmin, wal_status, safe_wal_size
4. inspect consumer, network, auth, schema and apply conflicts
5. calculate disk and slot-loss deadlines
6. choose resume, repair, rebootstrap or retire
7. verify acknowledgement plus semantic convergence
8. only then close or drop the slot

优先恢复 consumer,而不是先 drop slot。若 slot 已 lost,继续重试同一位点没有意义; 冻结下游写入,按 snapshot + new slot 的协议重建。若 consumer 已永久退役,保留审批与 最后消费位置后精确 drop。

与 Pigsty 观测对齐

Pigsty 的 PGSQL Replication / Persist / Instance / Alert 看板可把:

slot retention
WAL production
archive
disk free
physical replica lag
logical pub/sub
host I/O

放在同一时间轴。原生 SQL 则确认 slot、subscription、table state 与 conflict identity。 平台看板用于发现趋势,不能替代 consumer owner 和业务 reconciliation。

进一步阅读:


上一节:逻辑复制原语 · 返回本章目录 · 下一节:批量装载与数据校验 · 查看全书目录 · 查看索引中心

29.3 批量装载与数据校验

迁移最容易制造一种虚假的成功感:目标端已经有很多行,增量也在流动,于是团队宣布 “数据迁完了”。但全量装载只回答“怎样把字节搬过去”,数据校验才回答“搬过去的是否 还是同一份业务事实”。

本节把两者作为一个不可拆分的阶段:装载方案必须预先定义验证方法,验证失败必须能 定位到批次、分桶乃至具体主键,而不是在切流前夜才比较两个 count(*)

29.3.1 COPY、并行、约束和索引顺序

先区分三条全量路径

路径 一致性边界 适用场景 主要代价
subscription initial copy 由 table sync worker 与 slot 协调 PG 到 PG,目标表已准备好 并行度和变换能力受逻辑复制模型约束
pg_dump / pg_restore dump snapshot 完整或选择性对象迁移 需要自行衔接 dump 后的增量
COPY / \copy 管道 由导出事务和位点协议定义 大表、异构转换、分批装载 快照、分片、错误账本和增量汇合都要自己负责

COPY 很快,但它不自动提供迁移一致性。若导出事务没有与 logical slot 的 exported snapshot 对齐,逐表 COPY 得到的可能是不同时间点;若完成全量后才创建 slot,全量与 增量之间还会留下永久缺口。第 29.2 节的“snapshot 加 stream”协议因此同样适用于手工 批量装载。

PostgreSQL 中有两个经常混淆的文件边界:

COPY shop.orders (order_id, customer_id, status, amount, updated_at)
TO '/server/path/orders.csv'
WITH (FORMAT csv, HEADER true, ENCODING 'UTF8');

COPY 的文件由数据库服务器进程读取或写入,需要服务器文件权限;psql\copy 则让客户端读写文件,通过 SQL 连接传输数据。迁移工作站通常使用 \copy, 避免给数据库角色服务器文件权限。无论选哪一种,都应:

  • 显式列出列名,不依赖物理列顺序;
  • 固定编码、日期格式、时区和 NULL 表示;
  • 记录导出查询、snapshot、源系统标识、行数、文件大小与文件摘要;
  • 把原始文件或不可变对象版本作为可追溯输入;
  • pg_stat_progress_copy 观察正在执行的 COPY,而不是从文件大小猜完成度。

binary COPY 省去文本转换,在完全同构、版本和类型实现已验证时可能更快;它不是通用 交换格式。跨 major、跨架构或有类型映射时,文本/CSV 加显式规范通常更可审计。

并行单位要可重放

一条 COPY 不能通过加一个参数变成并行任务。常见并行单位是:

不同表
同一分区表的不同叶子分区
按稳定主键范围切片
预先生成且有 manifest 的多个文件
pg_restore 的独立对象任务

切片必须互斥、完备并可复算。例如按整数主键范围切分时,记录 [lower, upper),不要用随数据变化的 LIMIT/OFFSET。按 hash 分桶时,固定 hash 算法、编码和桶数。并发量同时受源端顺序读、网络、目标 WAL、磁盘、索引维护、 autovacuum、standby 重放与连接数约束;“有 32 核就开 32 个 COPY”不是容量模型。

可先用一小段代表性数据测量:

source export MB/s
network MB/s and retransmission
target heap MB/s
WAL bytes / loaded byte
standby replay lag
checkpoint pressure
CPU spent on conversion and indexes

再逐级增加 worker,找到吞吐开始变平、延迟或 WAL 开始恶化之前的并发点。

约束、触发器和索引的顺序是风险选择

COPY FROM 会执行 check constraint 和 trigger,但不会执行 rewrite rule;外键检查、 二级索引维护和触发器都可能成为装载成本。不能因此笼统地把它们全部关闭:

做法 收益 风险与前提
保留 PK/UNIQUE/CHECK 立即拒绝重复或非法行 装载时持续维护索引
装完再建二级索引 批量排序建索引通常更快 装载期间查询能力弱,建索引需额外空间
按父表再子表装载 可保留 FK 检查 并行度下降
装入 staging 再转换 错误隔离、类型转换可审计 多一份空间与一次写入
暂缓 FK 后再 VALIDATE 加快大批量导入 切流前必须完成验证,且不能让非法数据外泄

对 online migration,目标表通常已经服务 logical apply。随意禁用 trigger、 session_replication_role 或删除 replica identity,可能同时改变增量应用语义。正确 顺序应在演练中固化,例如:

创建 schema 与必要主键
  -> 创建不参与装载路径的必要类型/扩展
  -> 全量装载或启动 initial copy
  -> 建立可延后的二级索引
  -> 验证/启用约束
  -> ANALYZE
  -> 等待增量追平
  -> 运行数据与业务校验

装载后立即 ANALYZE。否则数据虽然完整,优化器仍可能按空表或旧统计量选择计划, 把“迁移正确”误判成“新库性能不行”。

29.3.2 行数、摘要、分桶与业务不变量

校验是一架逐层缩小范围的梯子

单独的 count(*) 很弱:删掉一行再插入一行,行数完全不变。反过来,直接对十亿行做 一个全表摘要虽然更强,一旦不一致却只会得到“某处不同”。实用校验从便宜到昂贵逐层 推进:

  1. 对象 manifest:schema、表、列、类型、默认值、identity、约束、索引、分区、 publication membership;
  2. 精确行数:不能拿 pg_class.reltuples 这类估算值做最终验收;
  3. 列统计min/max/sum/null count/distinct count、状态分布;
  4. 稳定有序摘要:对规范化后的逻辑行计算 digest;
  5. 分桶摘要:发现差异后只重扫异常桶;
  6. 业务不变量:外键孤儿、金额边界、状态机、账务守恒;
  7. 代表性业务查询:从应用可见结果验证语义与性能。

本章实验把一张表的 logical manifest 表示为:

row_count
ordered row digest
numeric sum where applicable
status histogram where applicable
invariant violations

源端和目标端都用同一组显式列与规范化规则生成它,而不是比较 heap 文件或物理 WAL。 初始复制的正式证据为:

customers = 5,000
orders    = 20,000
两张表 pg_subscription_rel 状态均为 r
源、目标 logical manifest 完全相同

后续又同步 500 个 insert、200 个 update 和 100 个 delete,等 marker 被目标确认后再次 比较,manifest 仍完全相同。

摘要必须先定义规范化

下面这种拼接并不可靠:

md5(string_agg(a || '|' || b, '' ORDER BY id))

因为 NULL、分隔符转义、浮点格式、timestamp 时区、JSON key 顺序、collation 和编码 都可能制造歧义。更安全的合同至少明确:

columns: [order_id, customer_id, status, amount, updated_at]
order_by: [order_id]
null_token: "\\N"
text_encoding: UTF-8
numeric_scale: 2
timestamp_zone: UTC
timestamp_precision: microseconds
json_canonicalization: sorted-keys
row_framing: length-prefixed
digest: sha256

摘要算法不是安全认证;它是高概率发现迁移差异的工程手段。关键业务金额还应比较精确 聚合和业务不变量,不能只依赖 hash。

分桶让差异可定位

以不可变主键把行分成固定数量的桶:

bucket = stable_hash(primary_key) mod 16

每个桶分别记录行数和摘要。正式实验在目标端只改动 order_id = 1,16 个桶中只有 bucket 1 不一致;从源权威行修复后,不一致桶集合回到空。这比发现全表摘要不同后重新 传输整张表更适合持续 reconciliation。

生产中可递归细分:

table mismatch
  -> bucket mismatch
      -> primary-key range mismatch
          -> row-level diff
              -> approved repair

修复操作也要写 ledger:源权威端、主键、修复前后摘要、执行者、ticket、commit time 和复核结果。不要让“校验工具”直接静默覆盖目标。

校验也会与写入竞态

如果源端仍在写,先扫源、再扫目标,结果可能来自不同逻辑时点。可选方案包括:

  • 在同一个 exported snapshot 上导出基线;
  • 记录源端 marker LSN,等待目标确认后再比较;
  • 对持续校验连续运行两轮,只升级稳定重复的差异;
  • 按业务 updated_at 水位排除仍在变化的尾部;
  • 在冻结窗口内做最终强校验。

“这次比较相等”必须附带比较边界。否则它只能证明两个扫描偶然读到了相同结果。

29.3.3 装载速度不能牺牲可追溯错误

PostgreSQL 18 的 COPY FROM 可以对文本或 CSV 输入使用:

COPY migration_stage.orders_raw
FROM STDIN
WITH (
  FORMAT csv,
  HEADER true,
  ON_ERROR ignore,
  REJECT_LIMIT 100,
  LOG_VERBOSITY verbose
);

这给“少量脏行继续装载”提供了原生工具,但边界很窄:

  • ON_ERROR ignore 只忽略把输入字段转换为目标类型时的错误;
  • constraint、trigger、I/O 等错误不会因此都被吞掉;
  • REJECT_LIMIT 应是显式且很小的错误预算,超过立即失败;
  • verbose 日志可能包含输入值,只能进入受控证据目录;
  • 被忽略的行必须进入后续补录与复核流程,不能只在日志中存在。

一个可追溯 reject 账本至少保存:

run_id: 2026-07-29-shop-orders-01
source_object: s3://migration/orders/part-017.csv
source_sha256: ...
record_locator: line-18342
primary_key_if_known: 923812
error_class: invalid_numeric
raw_record_ref: encrypted://...
decision: pending
repair_version: null
replay_run_id: null

原始敏感行不必进入普通日志;可以保存不可逆摘要和受控对象引用。重要的是能够回答: 这行来自哪里、为什么被拒绝、是否修复、在哪个 run 重放、最终是否进入目标。

staging 比在正式表里猜错更便宜

异构或质量未知的数据优先装入 staging:

raw text columns
  -> parse and classify
      -> quarantine rejects
          -> cast into typed staging
              -> validate business rules
                  -> merge into target

这样 conversion error、业务 error 和目标冲突可以分别统计。正式表上的 transaction 仍应保持全成或全败;批次间可独立提交,但每个批次必须有不可变输入和 idempotent 重放方法。

一次大 COPY 在中途失败并回滚后,已经插入的 tuple 会成为不可见 dead tuples,占用 空间,之后可能需要 VACUUM 回收。把重试理解为“失败就再跑一次”会在有限窗口里放大 I/O 和磁盘压力。应在演练中测量失败批次的空间后果,合理拆批,并为 vacuum 留预算。

本阶段的停止线

满足以下条件,才能从“全量装载”进入“增量追平”或最终校验:

  • 每个输入文件/切片都有 manifest、行数与摘要;
  • 成功行数加拒绝行数与输入记录数守恒;
  • reject 未超过预算,且每一行都有处置状态;
  • 目标对象、精确行数、分桶摘要和业务不变量已输出;
  • deferred index 已创建,constraint 已验证,统计信息已更新;
  • 任一失败批次都能无副作用重放;
  • 校验采用的 snapshot/marker 边界已记录。

吞吐是迁移的约束,不是迁移的正确性定义。一个快到无法解释丢了哪些行的装载流程, 不具备上线资格。


上一节:CDC 与复制槽治理 · 返回本章目录 · 下一节:在线迁移状态机 · 查看全书目录 · 查看索引中心

29.4 在线迁移状态机

在线迁移不是一条命令,而是一个有进入条件、退出证据和失败转移的状态机。把 runbook 写成“先全量,再增量,最后切流”,现场仍会争论:什么叫追平、何时禁止旧库写入、目标 写过以后还能否回退。

本节把这些模糊动词改成可观测状态。每次转移都要求证据;证据不足就留在原状态,不用 截止时间替代正确性判断。

29.4.1 预检查、全量、增量、追平与冻结窗口

先定义状态,再填命令

状态 进入条件 退出证据 失败时动作
PREFLIGHT 迁移合同、owner、窗口已批准 兼容性、容量、权限、网络和回退演练通过 修合同,不创建长寿命 slot
BASELINE snapshot/slot 边界已建立 全量 manifest、reject ledger、对象验证通过 停装载,保留输入与证据后重建目标
STREAMING 基线与增量无缝衔接 subscription/consumer 稳定,目标持续跟随 修 apply/connector,不切流
CATCHUP 进入切换前观察 marker 已在目标可见,lag 和 retained WAL 入预算 降写入、扩容或推迟窗口
FROZEN 写围栏已生效 旧端写测试失败,最终 marker 和强校验通过 解除围栏或转前滚修复
CUTOVER 路由变更已批准 新连接身份正确,目标写 canary 成功 按回退矩阵决策
OBSERVE 目标承接生产流量 SLO、数据、任务、slot、日志持续合格 回退或前滚
EXITED 观察窗口和退出条件满足 迁移签收、资产和凭据收尾完成 不再把旧源当即时回退方案

状态不允许跳转,例如没有 FROZEN -> final validation,不能从仍在双写的 STREAMING 直接宣布 CUTOVER

preflight 要检查迁移语义,不只检查端口

预检查至少覆盖:

source/target system identifier, version, encoding, locale, timezone
schema, types, extensions, functions, collations, generated/identity columns
table size, churn, large transactions, partition topology
primary key and replica identity
RLS, owners, grants, triggers, rules
DDL change policy
sequences and identity allocation
large objects and external object references
publication filters and subscription options
WAL generation, slot retention budget, disk headroom
network throughput, latency, TLS/HBA and credential lifetime
application compatibility, connection pool and prepared statements
backup, restore, rollback and reconciliation procedure

“目标能连通”只证明网络路径存在。比如目标缺少 collation、sequence 没有同步、源表无 replica identity,都可能在增量或切流阶段才暴露。

Pigsty 的迁移任务可以生成环境检查、schema、publication/subscription、进度、差异和 sequence 操作的上下文与脚本;它不能替业务确认 trigger 语义,也不知道应用路由和 外部副作用。生成脚本应进入评审和版本控制,不能把“生成成功”当作“迁移完成”。

“追平”要有业务可见 marker

单看:

SELECT pg_wal_lsn_diff(pg_current_wal_lsn(), confirmed_flush_lsn)
FROM pg_replication_slots
WHERE slot_name = 'pg36_shop_slot';

只能看到 publisher 与 consumer acknowledgement 的位置关系。它不必然证明目标业务 查询已经看见某一笔事务,也不覆盖下游索引、缓存和异步任务。

更可靠的追平协议是:

  1. 在源端业务表或专用控制表提交唯一 migration_marker
  2. 记录该事务的业务 ID、提交时间和附近 LSN;
  3. 等目标端通过普通应用路径读到 marker;
  4. 同时确认 subscription worker、table state、slot、错误统计和 retained WAL 正常;
  5. 在一段稳定窗口中重复,而不是只采一个瞬时零延迟。

正式实验同步完 500 个 insert、200 个 update、100 个 delete 后,等待 marker 在目标端 可见,再比较两端 manifest。这个证据比“延迟图降到 0”更接近切换要求。

冻结窗口要证明旧写入真的失败

冻结不等于在群里发一句“请勿写库”。写围栏可以来自:

  • 撤销专用 runtime role 的 DML 权限;
  • 将旧端业务入口切换为只读;
  • 应用 feature flag 阻止写请求;
  • 停止 scheduler、ETL、CDC 回写和运维脚本;
  • 对无法配合的 writer 建立数据库级拒绝规则。

然后用旧应用凭据执行负向 canary,确认 INSERT/UPDATE/DELETE 失败,同时需要的 SELECT 仍可用于核对。表 owner、superuser、SECURITY DEFINER 函数和绕过业务入口 的后台任务必须单独盘点;仅 revoke 普通角色无法约束这些路径。

29.4.2 影子读、双读、切流和观察

影子读与双读解决的是“应用是否认同”

数据库摘要相同,不代表应用行为相同。目标端可能因为 collation、timezone、扩展版本、 查询计划或 session 参数给出不同结果。切流前可以逐级放量:

方法 主结果来自 目标端副作用 适合发现
离线 replay 源端录制流量 禁止 SQL/类型/性能不兼容
影子读 源端 严格禁止 结果、错误码、延迟差异
采样双读 源端,后台比较目标 只读 长尾和真实参数差异
小比例真实读 目标端 应用正常读副作用需评估 连接池、缓存、SLO

“读”也可能有副作用:SELECT nextval(...)、advisory lock、临时表、审计函数、缓存 填充、SELECT ... FOR UPDATE 都不适合直接镜像。影子层必须有 SQL allowlist、超时、 并发限制和结果脱敏;不能为验证目标把源端峰值流量翻倍。

结果比较应按业务语义规范化:

unordered result -> sort by declared key
timestamp -> normalize to UTC and agreed precision
numeric -> agreed scale and rounding
JSON -> canonical form
expected nondeterminism -> exclude or compare distribution

同时比较错误类别、行数、关键字段、P50/P95/P99 与资源消耗。只比较 HTTP 200 会漏掉 返回空集、排序变化和悄悄截断。

切流要拆开连接地址与数据权威

应用可能经过:

DNS
  -> VIP / load balancer
      -> HAProxy service
          -> PgBouncer
              -> PostgreSQL primary

迁移 runbook 必须指出实际控制点及缓存时间。改变 DNS 不会自动清掉旧 PgBouncer 连接;改变 HAProxy backend 也不会让应用已持有的 session 消失。切换步骤通常包括:

  1. 记录旧 route generation、目标 endpoint 和回退 endpoint;
  2. 降低或确认 TTL,准备 health check;
  3. 在源端建立写围栏并做最终 marker/校验;
  4. 刷新 sequence,保证目标下一个值高于已迁移最大值;
  5. 修改唯一的权威路由控制点;
  6. drain/重建旧池,拒绝新连接进入源端写服务;
  7. 新连接查询 system_identifier、database、server address 与只读状态;
  8. 执行目标 canary 并从应用层回读;
  9. 进入限时观察,不立即拆源。

本章实验刻意不修改真实 Pigsty 路由,只在私有证据中模拟 source -> target -> source。这证明状态机与回退逻辑,不声称实验真的改过平台入口。 生产切流必须由应用或网络 owner 执行并给出实际路由证据。

观察窗口看四层信号

关键证据
应用 成功率、错误分类、业务转化、队列积压、关键任务
连接/路由 新旧连接数、endpoint 身份、池等待、事务/会话模式
PostgreSQL TPS、延迟、锁、WAL、checkpoint、autovacuum、复制/订阅错误
数据 canary、分桶摘要、业务不变量、target-only writes、reconciliation

应预先写出阈值和观察时间,例如“连续 60 分钟错误率不高于基线 + 0.1%,关键不变量 为零,异常桶为零”。现场再决定“看起来还行”无法形成一致决策。

29.4.3 回退点、前滚点与不可逆动作

回退不是把连接串改回去

切流前,源是唯一 writer,回退通常只是解除源围栏并放弃目标。目标开始承接写入后, 状态发生根本变化:

source last state S
target receives new writes T1..Tn
route points back to source

若没有 reverse CDC 或显式 reconciliation,源端不知道 T1..Tn,直接回路由会造成 已成功请求消失。应在切换前确定矩阵:

所在阶段 目标端是否有独占写 默认策略
baseline / streaming / frozen 安全取消,恢复源写
cutover,canary 尚未产生业务事实 快速回退
cutover,只有可识别 canary 是,可枚举 对账并回放 canary 后回退
observe,已有普通业务写 reverse sync/reconciliation,或前滚修复
已执行破坏性 schema/源退役 是且难逆 按灾难恢复或专门回迁方案处理

正式实验在目标写入唯一 canary,得到 sequence 值 900001。模拟回退前,证据包先识别 并把这条目标独占数据对账到源,再在源写入 rollback canary,最终两端 manifest 相同。 这个演示只覆盖“可枚举的一条目标写”,不能推导出任意生产写流都可这样回退。

sequence 是典型的切换缝隙

逻辑复制不复制 sequence 当前值。正式实验切流前看到目标 sequence 当前值仍为 1, 而源端最大 order_id 已到 900000。流程显式执行 setval,目标 canary 才得到 900001

sequence 处理要考虑:

  • is_called 语义;
  • identity column 背后的实际 sequence;
  • 多 writer 是否预留不重叠区间;
  • cached values 与仍存活的旧连接;
  • 目标独占值回退后会否冲突;
  • gap 是否被业务错误地当成连续性失败。

标出不可逆动作

以下动作不应与普通切流混在一个“一键脚本”:

删除 source cluster/database
删除 backup 或缩短 WAL 保留
drop publication/slot before exit criteria
执行 target-only destructive DDL
重新使用旧 sequence 区间
轮换后销毁唯一可用的旧凭据
让外部系统产生不可补偿副作用

对 schema 使用 expand/contract:先部署双方都理解的扩展形态,完成迁移与观察,再在独立 变更中收缩旧字段。这样问题出现时可以前滚修复,而不是在数据迁移、应用发布、破坏性 DDL 三者同时发生时赌一个总回退按钮。

每次决策至少记录:

state: OBSERVE
route_generation: 42
source_fence: active
target_only_write_boundary: marker-20260729-17
rollback_possible: conditional
required_reconciliation: target-writes-after-marker
decision_owner: migration-commander
evidence_bundle: /secure/migrations/run-...
next_decision_deadline: ...

回退能力是一项需要实验证明的属性。未演练、未定义数据合并边界的“随时可回退”,只是 一句安慰。


上一节:批量装载与数据校验 · 返回本章目录 · 下一节:异构同步的语义损失 · 查看全书目录 · 查看索引中心

29.5 异构同步的语义损失

异构同步可以让目标端“有数据”,却无法自动保证两边表达的是同一件事。connector 显示 running、offset 持续推进、目标查询也返回 200,只说明管道在工作;类型舍入、 排序规则、事务边界和删除语义仍可能已经变化。

本节给出一套语义合同。它不仅适用于 PostgreSQL 到 MySQL、Kafka、Elasticsearch 或 数据仓库,也适用于两个配置、扩展与 locale 不同的 PostgreSQL 环境。

29.5.1 类型、精度、排序规则与时区

类型映射必须是一张可测试的合同

不能只写:

numeric -> decimal
timestamp -> timestamp
jsonb -> json

至少要写清:

源语义 目标映射需要回答
numeric(p,s) 最大精度、scale、舍入模式、溢出是失败还是截断
bigint / unsigned integer 目标上下界,超界行如何隔离
real / double precision NaN、正负无穷、负零、比较语义
char(n) / text 尾随空格、Unicode normalization、空串与 NULL
timestamp without time zone 它代表本地墙上时间还是业务约定 UTC
timestamp with time zone 输出 zone、精度、DST 重叠/缺口
jsonb key 顺序、重复 key、numeric 精度、缺失与 JSON null
UUID / enum 原生类型还是 text,非法值和新增 enum label
array / range / multirange 展开、序列化还是目标原生类型
bytea 编码、大小上限、二进制是否被误当字符串

应为每一种映射准备 boundary corpus,而不是只测正常样本:

min/max
刚好超界
0 / -0
小数临界舍入
NULL / empty
非 ASCII 与组合字符
DST 切换前后
闰日
超长值
NaN / Infinity where supported

迁移前后都用同一个 canonical encoder 输出,比较规范化值和预期错误类别。若业务决定 允许损失,例如金额从 4 位小数舍入到 2 位,必须记录舍入规则、受影响行数、总误差和 批准人;不能让驱动默认转换替团队做决定。

时区问题常被样本掩盖

PostgreSQL 的 timestamptz 保存一个绝对时间点,显示受 session TimeZone 影响; timestamp 不含时区。把前者格式化为本地字符串再写进后者,会永久丢掉 offset。

合同应明确:

source_type: timestamptz
wire_form: RFC3339 with numeric offset
canonical_zone: UTC
precision: microseconds
target_type: timestamp(6) with time zone
ambiguous_local_time_policy: reject

还要验证 connector、JDBC/driver 与 sink session 的时区,而不只比较 database 参数。 夏令时地区的 02:30 可能不存在,01:30 可能出现两次;用七月的一条 UTC 样本无法 覆盖这些边界。

collation 会改变“同样查询”的结果

字符值逐字节相同,也可能因 libc/ICU/provider/version 不同而产生:

  • ORDER BY 顺序变化;
  • case/accent insensitive 比较变化;
  • UNIQUE index 对“相等”的判断不同;
  • prefix/range query 命中集合不同;
  • 分页边界漂移。

迁移 inventory 应记录数据库和列级 collation/provider/version,并在目标查询 pg_collation 与实际索引定义。若应用依赖稳定顺序,应在 SQL 中给出完整 tie-breaker, 例如 ORDER BY display_name COLLATE ..., customer_id;只靠隐含排序,本来就没有 跨环境保证。

本章正式实验使用 PostgreSQL 18.6 到 PostgreSQL 18.6,且两端都由同一 Pigsty 实验环境管理。它能证明同构 PG 逻辑复制与校验流程,不能证明上述异构类型和 collation 合同。异构结论必须在真实 source/sink 组合上另做边界语料实验。

29.5.2 约束、事务顺序与删除语义

源端约束不会自动变成下游约束

源端可以依赖:

PRIMARY KEY / UNIQUE
FOREIGN KEY
CHECK
EXCLUDE
domain constraint
trigger-maintained invariant
transaction isolation
deferred constraint

消息流通常只携带行变化,不携带这些证明。目标是搜索索引或对象存储时,甚至没有对应的 约束机制。于是“source 每次提交都合法”不能推出“sink 任意时刻都合法”。

例如源事务先创建 customer 再创建 order。若 connector 按 table 分 topic,下游并行 消费,order 可能先可见。解决方式不是祈祷消费者够快,而是明确:

  • 是否保留 source transaction ID 和 commit boundary;
  • 跨表事件是否要求原子可见;
  • 不要求原子时,查询层如何隐藏未完成 batch;
  • parent 缺失是重试、暂存、告警还是丢弃;
  • checkpoint 在整个事务之后还是每条事件之后推进。

事务内 row order 也不能随意打散。账户扣款、入账和 ledger 三条事件若被三个 worker 独立提交,中间态会破坏守恒。高吞吐设计必须说明它牺牲了什么可见性,以及如何恢复。

upsert 需要版本,delete 需要墓碑

一个简单的:

INSERT ... ON CONFLICT DO UPDATE

只保证当前语句不因 key 冲突失败,不保证旧事件不会覆盖新状态。目标记录通常需要 source version/commit position,并采用条件更新。

删除则至少有四种不同语义:

源动作 下游可能需要
physical DELETE key tombstone,删除投影
soft delete 保留记录并同步 deleted_at
FK cascade 每个子变化或可重建的级联合同
TRUNCATE 清空整个 collection,或明示不支持并触发重建

若 sink 先收到 DELETE,随后重放一条旧 UPDATE,没有 version/tombstone ledger 就会把 已删除对象复活。tombstone 的保留时间必须长于最大 replay/backfill 窗口;过早压缩会 重新暴露复活风险。

PostgreSQL publication 可以发布 TRUNCATE,但 row filter 不会过滤它。外部 CDC connector 是否把它转成一个控制事件、逐行 delete 还是直接不支持,要在上线前实测。

backfill 与实时流必须共享所有权规则

backfill 可能比实时事件更晚到:

snapshot contains version 7
stream has already applied version 9
backfill blindly upserts version 7

目标就回到了旧状态。每个写入路径都必须服从同一条条件:

apply only if incoming source version is newer
or if this batch is the declared authoritative rebuild

重建期间可以使用新的目标 namespace/index/table,完成校验后原子交换;不要让不带 version 的历史 backfill 与实时流争写同一记录。

29.5.3 目标端可查询不等于语义等价

绿灯只能证明它声明的那一层

绿灯 能证明 不能证明
connector running 进程存活并执行主循环 没有跳过 poison event
offset advancing 一些事件被确认 sink 副作用完整、顺序正确
target row count 相等 总行数一致 行内容、关联和删除一致
target query 成功 语法和服务可用 排序、精度、完整性等价
lag 接近零 消费接近 source head 历史基线正确

异构验收应沿一条更强的梯子:

transport alive
  -> no unaccounted rejects
      -> schema/type contract passes boundary corpus
          -> row and bucket manifests agree
              -> business invariants agree
                  -> representative queries agree
                      -> workload SLO agrees
                          -> reconciliation remains stable over time

代表性查询不是随机挑十条 SELECT *,而应从业务清单中覆盖:

  • equality、range、prefix、全文与排序;
  • NULL、缺失字段、数组/JSON 嵌套;
  • pagination 和 tie-breaker;
  • 聚合、去重、金额与时区窗口;
  • 删除、恢复、乱序和重复事件;
  • 最大对象、热点 key 与大事务;
  • 权限过滤和租户隔离。

每条都定义允许差异。例如搜索结果可能允许排名小幅变化,但不能跨租户;报表金额必须 精确相同;分析仓库允许 10 分钟最终一致,但 reconciliation 不允许永久缺口。

建立“允许损失登记表”

异构系统很少完全同构,现实做法不是假装零损失,而是让损失显式:

field: customer.display_name
difference: ICU collation produces different tie order
affected_queries: customer-search
business_impact: none when customer_id is secondary key
mitigation: append customer_id to ORDER BY
validation: query-corpus/collation-03
owner: customer-platform
approved_until: permanent

没有登记的差异一律视为 defect;登记项也要有 owner、验证和复审条件。这个机制防止 “已知差异”在口头交接中无限扩张。

权威源和修复方向必须唯一

持续 reconciliation 发现不一致时,先回答:

在当前阶段谁是 source of truth?
差异来自漏事件、重复、乱序、手工写还是映射改变?
修复目标会不会被下一条旧事件再次覆盖?
修复需要 rewind、rebootstrap 还是单 key replay?
该修复怎样留下 provenance?

切流前通常以源端为权威;切流后目标已承接新写,不能继续无条件“用源覆盖目标”。 权威边界随迁移状态改变,必须随状态机一同记录。

本章的目标端 drift 实验很能说明这一点:目标端手工修改 order_id = 1 后,subscription 仍是 running,目标查询也正常,但只有 bucket 1 的摘要暴露了差异。因为当时仍处于 切流前阶段,流程才能用源端权威行修复。若那是一笔切流后的合法目标写,相同动作反而 会销毁正确数据。

所以,数据“到了”是传输结论;业务“等价”是由类型合同、事务合同、校验语料和持续 对账共同支持的结论。两者不能用同一个绿色图标代替。


上一节:在线迁移状态机 · 返回本章目录 · 下一节:多集群迁移环境 · 查看全书目录 · 查看索引中心

29.6 多集群迁移环境

Pigsty 能把两个 PostgreSQL 集群、服务发现、连接池、监控和配置组织在同一个管理面里, 但迁移的数据面仍跨越两个独立系统。最危险的自动化错误,往往不是 SQL 写错,而是把 命令发到了正确环境中的错误 cluster、database 或 service。

本节把端点身份、凭据、观测和退出窗口做成多集群迁移的硬边界。

29.6.1 源、目标、验证端点与隔离凭据

给每类动作一个明确端点

端点 连接对象 允许动作 不应承担
source migration 源 primary/direct-write service 建 publication/slot、读快照、写 marker 普通应用写
source runtime 源业务 service 迁移前业务读写,冻结后负向 canary 管理 slot
target migration 目标 primary/direct-write service 建 schema/subscription、装载、sequence 源端管理
target runtime 目标业务 service 影子读、切流后业务读写 超级用户迁移 DDL
verification 两端只读连接 manifest、目录、业务不变量 修复或路由变更
monitoring exporter/API/dashboard 采集两端与代理指标 作为业务正确性的唯一证据

源 logical replication 连接必须到能够创建 slot、持续发送 WAL 的 primary。目标 subscription 在目标 database 中创建,apply 写目标表。不要把“read service”名称当作 复制端点,也不要把 HA service 的自动切换能力误认为 slot 已具备 failover 语义。

每次 run 先打印并保存非秘密身份:

SELECT current_database(),
       current_user,
       inet_server_addr(),
       inet_server_port(),
       current_setting('server_version'),
       pg_is_in_recovery();

SELECT system_identifier
FROM pg_control_system();

并断言 source 与 target system_identifier 不同。database 名称相同不等于同一个 数据源,IP 不同也不保证不是同一 cluster 的两个实例。

凭据按职责和环境隔离

至少拆分:

source runtime
target runtime
replication/login
migration DDL
verification read-only
monitoring

replication 角色需要源端 LOGIN REPLICATION,initial copy 还需要 published table 的 SELECT;它不需要目标端 DDL。目标 runtime 不应能回写源。验证角色不应因为要比较 数据而获得修复权限。

连接合同还包括:

  • TLS mode、CA 与 hostname verification;
  • 精确 HBA source CIDR 和 role classification;
  • CONNECT、schema USAGE、table privilege;
  • secret manager 引用、有效期和轮换 owner;
  • conninfo 禁止出现在公开证据、进程参数或 shell trace;
  • 迁移完成后的 revoke/rotate 清单。

本章第一次候选实验在插入 fixture 之前就被 Pigsty HBA 拒绝:临时登录角色没有加入 HBA 所使用的平台分类角色。实验没有放宽成全网 trust,而是把 runtime 与 replication 账号分别加入 dbrole_readwrite / dbrole_readonly,并使用 INHERIT FALSE, SET FALSE,只满足成员分类,不继承平台对象权限。失败候选随后按 marker 精确清理,确认两端 database、role、slot 与 subscription 都不存在。

这条失败很有价值:HBA 的“能否建立连接”和 SQL grant 的“连上后能做什么”是两道 独立门。

Pigsty 迁移上下文是编排模板,不是控制平面替身

Pigsty 的 pgsql-migration.yml 可以围绕 source/target inventory 生成并组织:

check-user / check-db / check-hba / check-repl / check-misc
copy-schema
create-pub / create-sub
copy-progress / copy-diff
copy-seq

以及需要 operator 落实的 source disable 和 re-routing 步骤。使用时应:

  1. 固定 Pigsty/模板版本;
  2. 审阅生成的变量、端点和 SQL;
  3. 把生成物与迁移 ticket 绑定;
  4. 在 disposable database 完整演练;
  5. 由实际应用/网络 owner 实施写围栏与切流;
  6. 用 PostgreSQL catalog 和应用身份反向复核。

本章正式实验在本地开发环境以 pg-test 为源、pg-meta 为目标,各自创建一次性 database 和角色。这只是为了在有限实验环境里证明真正的双 cluster 边界。生产中不应 因为教程这样做,就把承担 Pigsty 管理面的 pg-meta 当成默认业务迁移目标。

29.6.2 观察槽、WAL、延迟与切换流量

先用原生视图定义信号

源端:

SELECT slot_name,
       database,
       active,
       active_pid,
       restart_lsn,
       confirmed_flush_lsn,
       pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS retained_bytes,
       wal_status,
       safe_wal_size,
       inactive_since,
       invalidation_reason,
       failover,
       synced
FROM pg_replication_slots
WHERE slot_name = 'pg36_shop_slot';

目标端:

SELECT subname,
       worker_type,
       pid,
       received_lsn,
       latest_end_lsn,
       last_msg_send_time,
       last_msg_receipt_time,
       latest_end_time
FROM pg_stat_subscription
WHERE subname = 'pg36_shop_sub';

SELECT *
FROM pg_stat_subscription_stats
WHERE subname = 'pg36_shop_sub';

再加上 pg_subscription_rel table state、PostgreSQL log 和业务 marker。不同版本的 视图列会变化,自动化应先断言 server major,而不是对未知列 SELECT * 后按位置解析。

关键关系是:

consumer stopped
  -> confirmed position stops
      -> restart_lsn cannot advance
          -> retained WAL grows
              -> disk pressure or slot invalidation

正式实验禁用 subscription 后确认 slot inactive,生成 3,000 条变化:

retained WAL: 227,008 -> 2,867,128 bytes
confirmed_flush_lsn: unchanged during stall

重新启用后才追平。这个小规模实验不能给生产容量一个固定阈值,却证明“subscriber 不工作”会转化为 publisher 的磁盘风险。

max_slot_wal_keep_size 可以限制 checkpoint 时 slot 允许保留的 WAL;超过后 slot 可能失去继续消费所需的 WAL,而不是自动帮 consumer 修好。PostgreSQL 18 的 idle_replication_slot_timeout 可以在 checkpoint 时使长期闲置 slot 失效,synced slot 有例外。二者都是保险丝,不能代替 owner、告警、rebootstrap 方案和容量预算。

Pigsty 视图负责聚合,catalog 负责定案

Pigsty 的 PGSQL、Replication、Persist、Service 与 Proxy 等 dashboard 可以统一观察:

source WAL and replication
target TPS/latency/locks
disk and checkpoint
service health and connection distribution
pool/proxy sessions

dashboard 适合关联趋势和告警;切流门禁仍应把关键 catalog query、时间范围和截图/导出 写入证据包。监控抓取间隔可能错过瞬时状态,面板也不会知道 marker 是否已被应用查询 读到。

logical slot failover 需要一整组前提

PostgreSQL 18 支持 failover logical slot,但不能只设置 failover = true 就宣称源端 primary 可无缝切换。完整设计还涉及:

  • primary 上 slot 标记为 failover;
  • standby 启用 sync_replication_slots
  • standby 配置 primary_slot_name 和可用的 physical replication slot;
  • primary_conninfo 包含可连接到正确 database 的配置;
  • physical standby 向上游提供所需 feedback;
  • 对需要的同步边界使用 synchronized_standby_slots
  • 切换前确认目标 standby 上对应 slot 已 synced 且位置安全;
  • connector/subscriber 重新连接后的 endpoint 与 timeline 测试。

这些参数影响可用性、WAL 保留和复制反馈,必须在与生产同构的故障切换演练中验证。 未做这套设计时,源 primary failover 可能要求重新 bootstrap;迁移窗口要把它列为明确 故障分支。

路由观测与数据复制是两条链

publication/subscription 不知道业务连接走 DNS、HAProxy、PgBouncer 还是配置中心。 切流证据至少包括:

route control plane generation
resolved endpoint before/after
new connection system identifier
old/new pool active sessions
source runtime DML denial
target runtime canary
rollback endpoint and drain status

监控图中的目标 TPS 上升只是旁证。最强证据是使用真实 runtime 凭据建立一个新连接, 验证目标系统身份并完成可回读 canary。

29.6.3 保留源环境直到退出观察窗口

保留不是继续双主写入

切流后的源环境应进入受保护状态:

runtime writes fenced
normal reads limited to verification
backup and required WAL retained
schema changes frozen
credentials and network path controlled
monitoring remains active
no scheduled job silently resumes writes

这样它既能支持调查和条件回退,又不会继续产生与目标分叉的新事实。若业务需要 reverse replication,应把它作为独立的数据合同:方向、冲突、sequence、DDL、RPO 和停止点都要 明确,不能把两个方向各建一个 subscription 就叫 active-active。

退出窗口用条件,不只用日期

观察窗口可以有最短时长,但退出还应同时满足:

  • 所有应用、worker、ETL、BI 和管理入口已迁到目标;
  • 旧 service/pool 不再产生业务连接;
  • 目标 SLO 经历了至少一个代表性高峰和关键批任务;
  • manifest、分桶和业务不变量连续通过;
  • 目标独占写边界与备份恢复已验证;
  • 告警、值班、容量和灾备 runbook 已移交;
  • 迁移临时 slot/subscription/reject 已有保留或清理决策;
  • source 回退 SLA 已到期,并有 owner 签收;
  • 资产、CMDB、DNS、secret 与文档已更新。

如果月末关账是系统最关键路径,只观察一个低峰小时即使所有图都绿色,也没有覆盖真实 业务周期。

清理按依赖方向进行

正常 subscription 与远端 slot 仍关联且源端可达时,DROP SUBSCRIPTION 会尝试删除 远端 slot。若先断网络或删除源 database,目标端 drop 可能失败,源端还可能留下孤儿 slot。退出流程应:

  1. 记录 subscription、publication、slot、owner 和最终位置;
  2. 停止/确认 consumer 不再需要;
  3. 在两端都可达时正常删除 subscription;
  4. 在源端确认 main slot 与 table-sync slot 均不存在;
  5. 再删除 publication 和迁移临时角色/权限;
  6. 轮换生产凭据;
  7. 将 source decommission 作为独立、可审计且明确不可逆的变更。

不要使用“删除所有 inactive slot”“终止所有连接”或 DROP DATABASE ... WITH (FORCE) 作为通用收尾。正式实验的清理只匹配本 run 创建的 database、role、subscription 和 slot;没有终止无关 session,普通 database drop 即成功。

保留窗口结束后,旧源也不应无限期成为影子生产系统。无限保留会继续消耗备份、补丁、 监控与安全治理成本,还让团队误以为随时能回到一个早已过期的数据副本。退出签收就是 把“临时可回退”正式转化为“目标为唯一权威,恢复走新体系”。


上一节:异构同步的语义损失 · 返回本章目录 · 下一节:实战:迁移 pg36_shop · 查看全书目录 · 查看索引中心

29.7 实战:迁移 `pg36_shop`

前六节建立了 logical replication、CDC、全量校验、迁移状态机、异构语义与多集群边界。 本节把它们压缩到一个可重复的 Pigsty 双集群实验:

pg-test / pg36_shop_src
  -> publication + logical slot
      -> pg-meta / pg36_shop_dst
          -> subscription

实验不是“看见几行复制过去”就结束,而要主动制造 consumer 停滞、显式 apply conflict 和不会报错的 silent drift,再完成 sequence 校准、写围栏、模拟切流、条件回退与精确 清理。

29.7.1 完成全量加增量同步

先读边界,再运行脚本

本章只允许在已确认的 Pigsty 四节点开发沙箱执行:

source       pg-test / PostgreSQL 18
target       pg-meta / PostgreSQL 18
data         deterministic synthetic fixture
capture      L0 read-only preflight
exercise     L2 bounded two-cluster disposable fixture
production   forbidden
real route   never changed

完整合同在 lab-contract.md,机器可读的环境、对象、数量和验收条件在 requirements.json,迁移状态与允许动作在 migration-contract.json,拓扑图在 topology.mmd

实验只创建以下固定名称对象,并用随机 run_id marker 证明所有权:

source database       pg36_shop_src
target database       pg36_shop_dst
publication           pg36_shop_pub
slot                  pg36_shop_slot
subscription          pg36_shop_sub
five exact fixture roles

它明确禁止读取现有业务表、修改 Pigsty inventory、修改 Patroni/持久参数、改变真实 HAProxy/PgBouncer/DNS/VIP、终止无关连接和 force-drop。marker、对象名或连接范围有一项 不匹配,runner 都失败关闭。

先做纯静态合同检查:

static/labs/ch29/task.sh lint

完整实验会创建和删除两个一次性数据库,必须在已确认沙箱中指定一个新的私有证据目录:

export PG36_EVIDENCE_DIR="$(
  mktemp -d "${TMPDIR:-/tmp}/pg36-ch29.XXXXXX"
)"

static/labs/ch29/task.sh all

若希望分步审阅:

static/labs/ch29/task.sh capture
static/labs/ch29/task.sh exercise
static/labs/ch29/task.sh verify
static/labs/ch29/task.sh review

capture 在任何写入前检查:

  • 两端 service、cluster、PostgreSQL major、primary 身份;
  • 两个不同的 system identifier;
  • source wal_level = logical
  • 目标 database、role、slot 和 subscription 起点不存在;
  • 第 19、23、25、28 章上游证据存在且环境边界一致;
  • 实验源文件散列与随后执行的版本一致;
  • 远端临时目录和证据目录满足私有权限。

创建 schema 与逻辑复制对象

夹具包含:

shop.customers     5,000 rows
shop.orders       20,000 rows

两表都有稳定主键,orders.customer_id 引用 customer;状态、金额和更新时间具有固定 业务约束。runner 在源端创建 publication 和 logical slot,在目标端创建相同 schema 与 subscription,然后等待 pg_subscription_rel 中两张表都从同步状态进入 r

正式参考 run 证明:

source system id  7668025967696967004
target system id  7668025945980641675
customers         5,000
orders            20,000
tables ready      2
logical manifest  equal

system identifier 是本次沙箱证据,不是读者环境中的预期常量。验收的是“两端不同且分别 绑定已声明 cluster”,不是数字本身。

让初始快照与持续变更汇合

初始复制完成后,源端执行:

INSERT    500 orders
UPDATE    200 orders
DELETE    100 orders
COMMIT    one unique migration marker

流程等待目标读取 marker,再比较两端精确行数、有序摘要、金额合计、状态分布和业务 不变量。参考 run 的结果为:

inserted / updated / deleted  500 / 200 / 100
marker acknowledged           true
logical manifest equal        true

这同时证明了 initial copy 与增量可以汇合,以及验证是在声明的 marker 边界之后执行。 它不证明生产大表所需时间,也不覆盖迁移期间的 DDL;后者仍需独立编排。

连接失败也是 preflight 结果

开发中的第一次候选 run 在 fixture 写入前被 HBA 拒绝,因为临时角色未匹配 Pigsty 的 group-role 分类。流程没有临时放宽认证,而是:

  1. 停止实验;
  2. 按 marker 清理两个 database、五个角色、slot/subscription;
  3. 证明所有临时对象不存在;
  4. INHERIT FALSE, SET FALSE 的成员关系满足 HBA 分类;
  5. 重新从新的 run 和空证据目录开始。

生产演练也应如此:preflight 失败说明合同不成立,不能在原 run 上一边改权限一边继续, 否则最终证据无法说明实际执行了哪套安全边界。

29.7.2 注入消费者停滞与数据差异

反例一:consumer 停了,风险留在 source

runner 精确禁用 pg36_shop_sub,确认源端 pg36_shop_slot inactive,然后在源夹具中 生成固定 3,000 行变化。参考结果:

confirmed_flush_lsn unchanged      true
retained WAL before                227,008 bytes
retained WAL after               2,867,128 bytes
retained WAL grew                  true
caught up after re-enable          true

禁用动作只匹配本 run 的 subscription;负载行数固定,不改变 max_slot_wal_keep_size,也不制造无限 WAL。数字取决于 tuple、full-page image、 checkpoint 和版本,教学结论是方向:

consumer 不确认 -> slot 不能推进 -> source retained WAL 增长

重新启用并追平后必须再次校验 manifest,不能因 worker 恢复 running 就进入下一阶段。

反例二:显式冲突会停 apply

实验先在目标端插入:

order_id = 900000, target payload

再在源端提交同一 key、不同值。目标 apply 命中唯一键,PostgreSQL 18 的 subscription 统计出现 insert_exists

confl_insert_exists  0 -> 1
apply_error_count    0 -> 1

此时 slot 仍可能存在,subscription 也仍是一个 catalog 对象,但 apply 已无法越过 冲突事务。修复流程必须先证明:

冲突表与主键
源端权威值
目标端冲突值
错误计数和日志时间
允许采取的修复方向

实验删除精确的目标冲突 fixture 行,让源事务重放,再等待 marker 与 manifest 收敛。 生产环境不能把“删除目标所有冲突行后重试”写成通用脚本;不同冲突可能代表合法的 target-only write。

反例三:静默漂移不会停 apply

runner 在目标端直接修改 order_id = 1。该行之后没有新的源变化,因此:

subscription remains healthy
no apply conflict is raised
target query succeeds

但 16 个稳定 hash bucket 中,只有 bucket 1 的行数/摘要不一致。实验处于切流前, 合同指定 source 为权威,因而按主键读取源行、记录修复前后摘要、修复目标并重跑所有 分桶,结果:

mismatched buckets before  [1]
mismatched buckets after   []

这组反例区分了两种故障:

故障 apply 是否报错 主要发现方式
唯一键/缺行等显式 conflict 通常会 worker/log/pg_stat_subscription_stats
目标手工写、错误 backfill 等 silent drift 不一定 持续 manifest、分桶和业务不变量

运行状态与数据等价必须分别验收。

29.7.3 验证、切流、回退并输出迁移证据包

切流前修正 sequence 并建立写围栏

逻辑复制已经把 order_id = 900000 复制到目标,但 sequence 本身没有随 DML 推进:

target sequence before  1
source max order_id      900000

runner 按源端最大 ID 校准目标 sequence。随后目标 runtime canary 得到 900001,证明 不会立刻与已迁移主键碰撞。

源端则撤销 runtime 的 DML 能力,用同一凭据实际发起 INSERT

SQLSTATE       42501
INSERT         denied
UPDATE         denied
DELETE         denied
SELECT         retained

“执行过 revoke”不是证据,旧凭据的负向操作才是。生产还要盘点 owner、 SECURITY DEFINER、scheduler 和其他 writer。

只模拟路由,不碰真实平台

私有 route-history.json 记录:

source -> target -> source

它只是一台状态机的模拟输入。runner 会检查:

Pigsty inventory unchanged
Patroni configuration unchanged
real HAProxy/PgBouncer/DNS/VIP unchanged
actual_platform_route_changed = false

在模拟 target 阶段写入一条可识别 canary。回退前把这 1 条目标独占数据显式对账回源, 再切回 source 并写 rollback canary。最终参考结果:

customers                      5,000
orders                        23,402
logical manifest equal          true
orphan orders                      0
negative amounts                   0
invalid statuses                   0
source retained through rollback   true

这里证明的是“在目标独占写可枚举时,条件回退协议可执行”。真实业务流量已经写入目标后, 是否回退仍取决于 reverse sync、对账能力和不可补偿副作用。

证据包必须能反驳伪成功

私有证据目录包含:

preflight-evidence.json
remote/migration-evidence.json
remote/route-history.json
remote-cleanup.json
negative-report.json
validation-report.json
public-summary.json
review.txt
source file hashes

review 会检查证据权限、schema、交叉字段、源文件 hash、私密信息与 public/private 边界。validator 不只验证成功样本,还要求:

29 declared counterexamples rejected
19 live evidence mutants rejected
11 source files hash-bound

也就是说,篡改 system identifier、初始行数、marker、retained WAL、conflict counter、 bucket repair、sequence、写围栏、路由边界或清理结论,都不能继续得到 pass。公开参考 摘要在 migration-run.json;它不含密码、conninfo、主机密钥 或原始行。

完成实验后可以对同一私有证据包重复审计:

static/labs/ch29/task.sh verify
static/labs/ch29/task.sh review

但不能把另一个 run 的证据目录与当前源文件拼接使用。

清理也是验收阶段

正常删除目标 subscription 时,PostgreSQL 同时删除远端 main slot。runner 随后分别在 两端确认:

source database absent
target database absent
all fixture roles absent
source slot absent
target subscription absent
ordinary DROP used
force drop used = false
unrelated sessions terminated = 0
remote temp absent

任何一项不成立,实验都不算完成。尤其不能为了让 CI 变绿而终止所有连接或删除所有 inactive slot。

从沙箱证据到生产迁移票据

本实验的最终决策仍是:

production_ch29_gate = pending

正式票据至少还要补:

  • 真实 schema/type/DDL/sequence/large-object inventory;
  • 数据量、写入峰值、全量时长与 slot WAL 容量压测;
  • source primary failover 和 consumer restart 演练;
  • TLS、HBA、secret rotation 与权限评审;
  • 应用影子读、连接池 drain、真实路由 owner 和变更窗口;
  • 业务不变量、允许语义损失与 reconciliation owner;
  • target-only write 边界、回退/前滚矩阵、不可逆动作;
  • 备份恢复、RPO/RTO、观察窗口和退出条件。

读者完成本节后,应能提交的不是一句“逻辑复制已同步”,而是一份可以回答同步了什么、 在哪个边界相等、故障怎样暴露、切流由谁执行、何时还能回退、如何证明已清理的迁移 证据包。


上一节:多集群迁移环境 · 返回本章目录 · 下一章:推陈出新:版本升级与回滚策略 · 查看全书目录 · 查看索引中心