移花接木:逻辑复制、迁移与异构同步
29 移花接木:逻辑复制、迁移与异构同步
逻辑复制能让数据流动,但“数据正在流动”离“迁移已经成功”还很远。
publication 不复制 DDL,sequence 不随表中 identity 值推进,大对象不在复制范围内; subscription 停止后,slot 仍可能继续保留 WAL 与 catalog rows;本地写入 subscriber 既可能造成会报错的 apply conflict,也可能留下完全不报错的静默漂移;在切换窗口里, 源端是否真的停止写、目标序列是否安全、目标新增写如何反向带回,决定了回退是不是一句 空话。
本章把迁移建模为一个有证据、有门禁、可失败关闭的状态机:
其中每一条箭头都有前置条件、观测、超时、停手方式与恢复路径。不能用最终行数相同跳过 中间状态,也不能把“DNS 已改”当成目标端数据、权限和业务语义已经可用。
本章目标
理解逻辑复制、CDC 和数据搬迁的状态机,用校验与可回退切换完成迁移,而不是把“数据能流动”误当成迁移成功。
读完本章,读者应该能够:
- 从 WAL、output plugin、publication、slot、subscription 和 apply worker 解释原生 逻辑复制的数据路径;
- 区分 initial table synchronization、持续 DML 与独立的 schema/sequence 迁移;
- 为表选择 primary key、
REPLICA IDENTITY USING INDEX或受约束的FULL; - 用
pg_replication_slots、pg_stat_subscription、pg_stat_subscription_stats与pg_subscription_rel定位停滞和冲突; - 解释 slot 为什么提供可恢复位点,却不能天然提供 exactly-once 的外部副作用;
- 为 CDC sink 设计 event identity、幂等提交、位点原子性与 replay 策略;
- 设计
COPY/dump/restore 的并行装载顺序,并保留 rejected rows; - 用行数、摘要、分桶、不变量和抽样分别验证结构与数据;
- 把在线迁移拆成 preflight、全量、增量、追平、写围栏、切流、观察和退出;
- 区分 rollback、forward repair 与已经跨过不可逆点的前滚;
- 为异构系统显式记录类型、精度、时区、排序、约束、顺序和 delete 语义损失;
- 用 Pigsty 生成迁移上下文、观察两端集群,但不把生成脚本误认为自动切流授权;
- 建立源、目标、验证和管理端点的独立凭据与网络边界;
- 在源端主库切换时评审 logical slot failover,而不是默认 subscription 会自动跟随;
- 完成一份含原始证据、公开摘要、回退记录和双端精确清理的迁移证据包。
一张图看清复制与迁移
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 沙箱集群之间执行:
正式 run:
初始复制:
| 对象 | 行数 | table sync state | logical manifest |
|---|---|---|---|
shop.customers |
5,000 | r |
相等 |
shop.orders |
20,000 | r |
相等 |
随后增量执行:
暂停 subscription 后写入 3,000 行:
| 证据 | 暂停前 | 暂停后 |
|---|---|---|
| slot active | false | false |
confirmed_flush_lsn |
不变 | 不变 |
| retained bytes | 227,008 | 2,867,128 |
恢复 subscription 后,marker 被确认且双端重新相等。
冲突与静默漂移:
切换与回退:
最后普通删除两端数据库、五个角色、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 重演完整双集群实验,阅读私有证据与公开白名单,练习如何从一个失败阶段 安全恢复,而不是只观察成功路径。
推荐学习路线
应用开发者:
重点是 schema/sequence 边界、数据校验、应用兼容、读写切换与异构语义。
平台工程师:
重点是 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 也不能因为出现在模板里
就进入仓库或公开证据。
核心资料:
- PostgreSQL 18:Logical Replication
- PostgreSQL 18:Publication
- PostgreSQL 18:Subscription
- PostgreSQL 18:Conflicts
- PostgreSQL 18:Monitoring
- PostgreSQL 18:Logical Decoding Concepts
- PostgreSQL 18:
pg_replication_slots - PostgreSQL 18:
COPY - Pigsty:Data Migration
- Pigsty:PGSQL Replication Dashboard
实验文件
执行:
all 会按相同顺序执行。exercise 会在两个沙箱集群上真实创建 subscription/slot、
生成 WAL 并注入冲突,只能在明确的一次性开发/测试环境运行。
本章验收
只有当迁移证据包能回答以下问题,才算完成:
“两端 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 中的 变更集合:
一个表可以属于多个 publication,一个 publication 可以有多个 subscriber。它可以限定:
但需要记住几个条件:
- 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:
正常情况下,这条命令同时:
- 在目标 catalog 创建 subscription;
- 连接源端;
- 在源端创建持久 logical slot;
- 启动 apply worker;
- 为待同步表启动受资源上限约束的 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 的变化流”的持久状态:
关键字段:
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 消费它。
生命周期必须成对
常规生命周期:
若目标端必须删除 subscription,而源端暂时不可达,不能假设远端 slot 也消失了。 PostgreSQL 允许先解除关联:
随后必须在源端按名称、database、plugin、active 状态精确检查并删除孤儿 slot:
这不是通用“清理所有 inactive slot”脚本。一个 inactive slot 可能只是计划内停机的 消费者;只有迁移 owner、保留策略和恢复点都确认不再需要时才能删除。
最小权限不是一个万能复制账号
源端连接角色至少需要:
若该角色不是 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 分类业务
连接。临时账号被授予:
这只让 HBA 的角色成员测试匹配,不让登录角色继承或 SET ROLE 到平台角色。把网络
分类与对象权限拆开,既能连接,也不把“能通过 HBA”误写成“自动拥有业务表权限”。
29.1.2 初始同步、流式变更与复制身份
初始快照与主 apply 流怎样汇合
逻辑复制的主路径是:
已有数据不能从未来的 WAL 中凭空恢复,所以每张表还经历:
状态在 pg_subscription_rel:
| code | 含义 |
|---|---|
i |
initialize |
d |
copying data |
f |
table copy finished |
s |
synchronized |
r |
ready,进入正常复制 |
查询:
本章正式 run 的第一次采样看到两张表分别处于 s 与 d,约 0.55 秒后才都进入
r。如果只看 subscription 已存在或 apply worker 有 PID,会过早宣布 initial copy
完成。
r 也只证明 PostgreSQL 的 table sync 状态。正式验收还比较:
结果为 5,000 customers、20,000 orders,双端 logical manifest 相等。
publication 动作与 initial copy 是两套选择
假设:
publish = 'insert' 只限制后续 DML。默认 copy_data = true 时,已有行仍会被初始复制。
因此不能用 publication operation list 推断目标基线只含某类事件。
row filter 与 column list 有各自的版本规则;跨版本迁移必须按源、目标最低版本检查。 多个 publication 若以不同 column list 重叠发布同一张表,并不是一个可随意叠加的投影 系统。变更 publication 后,还需要:
新表才进入 subscription catalog;是否 copy_data 要显式决定。REFRESH 不是 DDL
迁移,也不会为目标创建缺失表。
REPLICA IDENTITY 回答“改哪一行”
INSERT 只需把新值写入目标;UPDATE 和 DELETE 必须携带足以在目标定位旧行的身份。
默认身份是 primary key:
选择顺序通常是:
使用另一个索引:
显式 identity index 仍必须满足 unique、immediate、非 partial、列非空等约束。另一条
容易混淆的 PostgreSQL 18 能力是:源端使用 FULL 时,目标端可用符合条件的 B-tree
或 hash 候选索引辅助查找;候选不能是 partial,左侧首字段必须是表列而非表达式。
源端 identity 不是 FULL 时,目标端也必须有由相同或更少列组成的可用 identity。
没有合适 key 时:
这会发送整个旧行。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 拓扑本身就是数据模型的一部分,应进入设计评审。
监控正在发生什么
目标端:
源端:
目标显示 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 匹配目标:
不会自动创建 schema、table、type、extension、function、constraint、index、owner、 privilege 或 RLS policy。初始 schema 可用受版本控制的 migration 或:
但不能把这条管道当成无需审查的最终方案。pg_dump 输出中可能包含 extension、
owner、tablespace、security label、event trigger 依赖;源目标 major 不同时,还需要用
目标版本工具与官方兼容路径评审。
持续 DDL 应按兼容顺序编排。例如增加一个源端马上会写入的新列:
如果先改源,新的 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 仍为:
如果此刻允许目标接收默认 identity 写入,第一笔就可能生成已经存在的 key。正确顺序是:
示意:
本章先把目标 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 是表列,可以复制;两者不要混淆。还应逐项盘点:
这些对象需要独立迁移和验收,不能因为两张业务表在复制就默认“整个 database 已搬完”。
subscriber 不是自动只读副本
subscriber 是普通 PostgreSQL database。应用若能在复制目标表本地写入,会产生两类 后果。
第一类会让 apply 报错:
错误型 conflict 会停止复制,必须修复目标数据/权限,或在理解数据损失后显式跳过远端 事务。跳过不是“重试”,而是声明这笔源端事务不再应用到目标。
第二类可能不停止:
PostgreSQL 18 在 pg_stat_subscription_stats 记录多类 conflict;部分 origin-differs
信息依赖 subscriber 的 track_commit_timestamp。但没有报错或 counter 为零,仍不
能证明双端相等。
本章正式注入:
结果:
删除精确的目标冲突行后,worker 重放同一源事务并收敛。随后又只改目标
order_id=1,这一次没有等待 apply 报错,而是 16 桶摘要中的 bucket 1 不一致。
按源端权威行修复后,mismatch 才归零。
这说明迁移需要两条独立告警线:
只看其中一条,都会漏掉真实故障。
trigger 与 apply 语义
持续 apply worker 的 session_replication_role 为 replica,普通 origin trigger/rule
默认不执行;可显式启用 replica/always trigger。初始同步更像 COPY,会触发行级和
语句级 INSERT trigger。于是同一张目标表在 initial copy 与 steady state 可能经过不同
trigger 路径。
迁移前应盘点:
不要靠 trigger 在目标重新制造本应从源复制的副作用,也不要让 initial copy 重发邮件、 支付请求或 webhook。外部副作用必须有专门的 replay/idempotency 合同。
进一步阅读:
- PostgreSQL 18:Logical Replication
- PostgreSQL 18:Publication
- PostgreSQL 18:Subscription
- PostgreSQL 18:Restrictions
- PostgreSQL 18:Conflicts
- PostgreSQL 18:Security
- PostgreSQL 18:
pg_subscription_rel
返回本章目录 · 下一节: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:
output plugin 决定编码,不决定下游业务语义:
| plugin / path | 典型用途 | 重要边界 |
|---|---|---|
pgoutput |
内建 publication/subscription | 二进制 replication protocol,不是给人直接读的文本 |
test_decoding |
测试和理解 logical decoding | 测试插件,不是生产业务协议 |
| 第三方 JSON/connector plugin | 通用 CDC | 版本、扩展、schema envelope、升级兼容需独立负责 |
不要把“输出是 JSON”误认为 schema compatibility 已解决。事件仍需定义:
缺少 before image 时,DELETE 怎样定位?TOAST 列未变化时,connector 如何表示?列 rename 是新字段还是同一字段?这些都不是 JSON 格式本身能回答的。
SQL 解码 API 与 streaming protocol
教学或诊断可用:
peek 不消费,适合观察;get 会推进消费位置。对 pgoutput 应使用 binary API 或
replication protocol,而不是强行调用文本函数。
生产 connector 通常通过 replication protocol:
“客户端已经读到”“写进本地 buffer”“写到外部系统”“外部事务提交”“向 PostgreSQL 确认”是五个不同瞬间。connector 必须声明哪个瞬间触发 acknowledgement。
snapshot 与 stream 必须无缝衔接
新建 logical slot 时可导出一个 snapshot。这个 snapshot 表示:
于是通用 bootstrap 可以:
- 创建 slot 并取得 exported snapshot;
- 在另一个事务
SET TRANSACTION SNAPSHOT; - 全量导出 snapshot 中的基线;
- 从 slot 持续消费 B 之后的变化;
- 在下游把全量与增量合并。
内建 subscription 将这套协调封装在 table sync workers 中。自行开发 CDC 时,如果先
普通 SELECT 全表、很久以后才创建 slot,会在两者间漏变化;若先创建 slot却不消费也
不预算磁盘,会让 WAL 无界增长。bootstrap 顺序必须成为协议,不是运气。
四个位置不要混为一个“offset”
一个成熟 CDC 管道至少有:
它们应满足与具体协议一致的单调关系。最危险的情况是 connector 先确认 source,再把 buffer 异步写 sink:
反过来,sink 已提交而 checkpoint 尚未推进会造成 replay;这通常可以用幂等处理,永久 缺口却无法凭空修复。因此多数 CDC 更愿意接受 at-least-once,而不是冒险提前确认。
plugin、slot 与版本都要进 inventory
建议为每个 consumer 记录:
slot 本身不知道 consumer owner、SLO 或 sink 状态。没有外部 inventory,inactive slot 只能告诉你“现在没人连”,不能告诉你“可以删”。
29.2.2 至少一次、重复事件与下游幂等
PostgreSQL 已经明确允许 replay
logical slot 是 crash-safe 的,但它的当前位置只在 checkpoint 时持久化。服务器崩溃后, slot 可能回到较早 LSN,于是最近变化再次发送。网络断开、consumer 在 sink commit 后但 ack 前崩溃,也会产生相同结果。
因此正确假设是:
“测试十次没重复”不能升级为 exactly-once 保证。
exactly-once 是端到端属性
若 sink 也是 PostgreSQL,可以在同一个目标事务里同时写业务状态与 dedup ledger:
随后才向 source 确认。若事务已提交但 ack 丢失,replay 会命中唯一键,不重复副作用。
但如果副作用跨多个系统:
没有一个 PostgreSQL transaction 能原子覆盖全部。应使用:
- transactional outbox;
- sink 自身的幂等 key / conditional write;
- 去重 ledger;
- 可重建投影;
- saga/补偿;
- 对不可重复副作用的业务级 request identity。
不要用“处理成功后更新 offset”一句话掩盖多个 commit 之间的崩溃窗口。
事件身份怎么设计
只用 txid 不够:32-bit XID 会回卷,跨 cluster/database 也会重复。只用表主键也不够:
同一行可以变化很多次。一个实用 envelope 可包含:
具体 connector 能提供哪些字段取决于协议。核心要求是:
不要把消费时间戳或随机 UUID 当 event id;replay 时它们会变。
upsert 不自动等于幂等
下面的 sink 写法:
只有在“最后写入覆盖即可”且顺序严格时才近似幂等。它处理不了:
- 较旧事件晚到,覆盖较新状态;
amount = amount + delta被重复执行;- DELETE/tombstone 后旧 UPDATE 复活;
- 多个 source 同写一个 key;
- 事件 schema 变化导致部分字段保留旧值;
- 外部副作用已经发生。
更安全的投影常带 source version:
LSN 只在同一 source timeline/合同内有序;跨 source 合并仍需业务 version/vector 或冲突 规则。
transaction boundary 不能随意打散
源事务:
若 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 的事件:
都不是完整策略。应记录:
dead-letter queue 只是隔离区,不是数据正确性的垃圾桶。任何 approved skip 都要进入 reconciliation,并记录业务影响。
29.2.3 槽停滞、WAL 保留与磁盘风险
inactive 不等于无害
slot 与连接生命周期独立。consumer 下线后:
基础查询:
pg_wal_lsn_diff(current, restart_lsn) 是按当前时刻估算该 slot 的 WAL 保留距离,不等于
磁盘上所有 WAL 文件恰好这么大;checkpoint、archive、其他 slot 与 segment 粒度都会
影响实际 pg_wal。
把风险换成时间
若近期 WAL 产生率为 bytes/s,slot 当前保留 ,可用于增长的安全空间为 ,则最粗略的时间预算:
若配置 max_slot_wal_keep_size = M,距离 slot 可能在 checkpoint 后失去所需 WAL 的
预算:
实际告警应使用变化率和低水位:
只告警 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 增量收敛,再执行:
确认目标 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。
恢复:
验收同时要求:
“slot active 又变 true”仍不够。
停滞处置顺序
优先恢复 consumer,而不是先 drop slot。若 slot 已 lost,继续重试同一位点没有意义; 冻结下游写入,按 snapshot + new slot 的协议重建。若 consumer 已永久退役,保留审批与 最后消费位置后精确 drop。
与 Pigsty 观测对齐
Pigsty 的 PGSQL Replication / Persist / Instance / Alert 看板可把:
放在同一时间轴。原生 SQL 则确认 slot、subscription、table state 与 conflict identity。 平台看板用于发现趋势,不能替代 consumer owner 和业务 reconciliation。
进一步阅读:
- PostgreSQL 18:Logical Decoding Concepts
- PostgreSQL 18:Replication Settings
- PostgreSQL 18:
pg_replication_slots - PostgreSQL 18:Logical Replication Monitoring
- Pigsty:PGSQL Dashboards
上一节:逻辑复制原语 · 返回本章目录 · 下一节:批量装载与数据校验 · 查看全书目录 · 查看索引中心
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 的文件由数据库服务器进程读取或写入,需要服务器文件权限;psql 的
\copy 则让客户端读写文件,通过 SQL 连接传输数据。迁移工作站通常使用 \copy,
避免给数据库角色服务器文件权限。无论选哪一种,都应:
- 显式列出列名,不依赖物理列顺序;
- 固定编码、日期格式、时区和
NULL表示; - 记录导出查询、snapshot、源系统标识、行数、文件大小与文件摘要;
- 把原始文件或不可变对象版本作为可追溯输入;
- 用
pg_stat_progress_copy观察正在执行的COPY,而不是从文件大小猜完成度。
binary COPY 省去文本转换,在完全同构、版本和类型实现已验证时可能更快;它不是通用
交换格式。跨 major、跨架构或有类型映射时,文本/CSV 加显式规范通常更可审计。
并行单位要可重放
一条 COPY 不能通过加一个参数变成并行任务。常见并行单位是:
切片必须互斥、完备并可复算。例如按整数主键范围切分时,记录
[lower, upper),不要用随数据变化的 LIMIT/OFFSET。按 hash 分桶时,固定 hash
算法、编码和桶数。并发量同时受源端顺序读、网络、目标 WAL、磁盘、索引维护、
autovacuum、standby 重放与连接数约束;“有 32 核就开 32 个 COPY”不是容量模型。
可先用一小段代表性数据测量:
再逐级增加 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,可能同时改变增量应用语义。正确
顺序应在演练中固化,例如:
装载后立即 ANALYZE。否则数据虽然完整,优化器仍可能按空表或旧统计量选择计划,
把“迁移正确”误判成“新库性能不行”。
29.3.2 行数、摘要、分桶与业务不变量
校验是一架逐层缩小范围的梯子
单独的 count(*) 很弱:删掉一行再插入一行,行数完全不变。反过来,直接对十亿行做
一个全表摘要虽然更强,一旦不一致却只会得到“某处不同”。实用校验从便宜到昂贵逐层
推进:
- 对象 manifest:schema、表、列、类型、默认值、identity、约束、索引、分区、 publication membership;
- 精确行数:不能拿
pg_class.reltuples这类估算值做最终验收; - 列统计:
min/max/sum/null count/distinct count、状态分布; - 稳定有序摘要:对规范化后的逻辑行计算 digest;
- 分桶摘要:发现差异后只重扫异常桶;
- 业务不变量:外键孤儿、金额边界、状态机、账务守恒;
- 代表性业务查询:从应用可见结果验证语义与性能。
本章实验把一张表的 logical manifest 表示为:
源端和目标端都用同一组显式列与规范化规则生成它,而不是比较 heap 文件或物理 WAL。 初始复制的正式证据为:
后续又同步 500 个 insert、200 个 update 和 100 个 delete,等 marker 被目标确认后再次 比较,manifest 仍完全相同。
摘要必须先定义规范化
下面这种拼接并不可靠:
因为 NULL、分隔符转义、浮点格式、timestamp 时区、JSON key 顺序、collation 和编码
都可能制造歧义。更安全的合同至少明确:
摘要算法不是安全认证;它是高概率发现迁移差异的工程手段。关键业务金额还应比较精确 聚合和业务不变量,不能只依赖 hash。
分桶让差异可定位
以不可变主键把行分成固定数量的桶:
每个桶分别记录行数和摘要。正式实验在目标端只改动 order_id = 1,16 个桶中只有
bucket 1 不一致;从源权威行修复后,不一致桶集合回到空。这比发现全表摘要不同后重新
传输整张表更适合持续 reconciliation。
生产中可递归细分:
修复操作也要写 ledger:源权威端、主键、修复前后摘要、执行者、ticket、commit time 和复核结果。不要让“校验工具”直接静默覆盖目标。
校验也会与写入竞态
如果源端仍在写,先扫源、再扫目标,结果可能来自不同逻辑时点。可选方案包括:
- 在同一个 exported snapshot 上导出基线;
- 记录源端 marker LSN,等待目标确认后再比较;
- 对持续校验连续运行两轮,只升级稳定重复的差异;
- 按业务
updated_at水位排除仍在变化的尾部; - 在冻结窗口内做最终强校验。
“这次比较相等”必须附带比较边界。否则它只能证明两个扫描偶然读到了相同结果。
29.3.3 装载速度不能牺牲可追溯错误
PostgreSQL 18 的 COPY FROM 可以对文本或 CSV 输入使用:
这给“少量脏行继续装载”提供了原生工具,但边界很窄:
ON_ERROR ignore只忽略把输入字段转换为目标类型时的错误;- constraint、trigger、I/O 等错误不会因此都被吞掉;
REJECT_LIMIT应是显式且很小的错误预算,超过立即失败;- verbose 日志可能包含输入值,只能进入受控证据目录;
- 被忽略的行必须进入后续补录与复核流程,不能只在日志中存在。
一个可追溯 reject 账本至少保存:
原始敏感行不必进入普通日志;可以保存不可逆摘要和受控对象引用。重要的是能够回答: 这行来自哪里、为什么被拒绝、是否修复、在哪个 run 重放、最终是否进入目标。
staging 比在正式表里猜错更便宜
异构或质量未知的数据优先装入 staging:
这样 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 要检查迁移语义,不只检查端口
预检查至少覆盖:
“目标能连通”只证明网络路径存在。比如目标缺少 collation、sequence 没有同步、源表无 replica identity,都可能在增量或切流阶段才暴露。
Pigsty 的迁移任务可以生成环境检查、schema、publication/subscription、进度、差异和 sequence 操作的上下文与脚本;它不能替业务确认 trigger 语义,也不知道应用路由和 外部副作用。生成脚本应进入评审和版本控制,不能把“生成成功”当作“迁移完成”。
“追平”要有业务可见 marker
单看:
只能看到 publisher 与 consumer acknowledgement 的位置关系。它不必然证明目标业务 查询已经看见某一笔事务,也不覆盖下游索引、缓存和异步任务。
更可靠的追平协议是:
- 在源端业务表或专用控制表提交唯一
migration_marker; - 记录该事务的业务 ID、提交时间和附近 LSN;
- 等目标端通过普通应用路径读到 marker;
- 同时确认 subscription worker、table state、slot、错误统计和 retained WAL 正常;
- 在一段稳定窗口中重复,而不是只采一个瞬时零延迟。
正式实验同步完 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、超时、
并发限制和结果脱敏;不能为验证目标把源端峰值流量翻倍。
结果比较应按业务语义规范化:
同时比较错误类别、行数、关键字段、P50/P95/P99 与资源消耗。只比较 HTTP 200 会漏掉 返回空集、排序变化和悄悄截断。
切流要拆开连接地址与数据权威
应用可能经过:
迁移 runbook 必须指出实际控制点及缓存时间。改变 DNS 不会自动清掉旧 PgBouncer 连接;改变 HAProxy backend 也不会让应用已持有的 session 消失。切换步骤通常包括:
- 记录旧 route generation、目标 endpoint 和回退 endpoint;
- 降低或确认 TTL,准备 health check;
- 在源端建立写围栏并做最终 marker/校验;
- 刷新 sequence,保证目标下一个值高于已迁移最大值;
- 修改唯一的权威路由控制点;
- drain/重建旧池,拒绝新连接进入源端写服务;
- 新连接查询
system_identifier、database、server address 与只读状态; - 执行目标 canary 并从应用层回读;
- 进入限时观察,不立即拆源。
本章实验刻意不修改真实 Pigsty 路由,只在私有证据中模拟
source -> target -> source。这证明状态机与回退逻辑,不声称实验真的改过平台入口。
生产切流必须由应用或网络 owner 执行并给出实际路由证据。
观察窗口看四层信号
| 层 | 关键证据 |
|---|---|
| 应用 | 成功率、错误分类、业务转化、队列积压、关键任务 |
| 连接/路由 | 新旧连接数、endpoint 身份、池等待、事务/会话模式 |
| PostgreSQL | TPS、延迟、锁、WAL、checkpoint、autovacuum、复制/订阅错误 |
| 数据 | canary、分桶摘要、业务不变量、target-only writes、reconciliation |
应预先写出阈值和观察时间,例如“连续 60 分钟错误率不高于基线 + 0.1%,关键不变量 为零,异常桶为零”。现场再决定“看起来还行”无法形成一致决策。
29.4.3 回退点、前滚点与不可逆动作
回退不是把连接串改回去
切流前,源是唯一 writer,回退通常只是解除源围栏并放弃目标。目标开始承接写入后, 状态发生根本变化:
若没有 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 是否被业务错误地当成连续性失败。
标出不可逆动作
以下动作不应与普通切流混在一个“一键脚本”:
对 schema 使用 expand/contract:先部署双方都理解的扩展形态,完成迁移与观察,再在独立 变更中收缩旧字段。这样问题出现时可以前滚修复,而不是在数据迁移、应用发布、破坏性 DDL 三者同时发生时赌一个总回退按钮。
每次决策至少记录:
回退能力是一项需要实验证明的属性。未演练、未定义数据合并边界的“随时可回退”,只是 一句安慰。
上一节:批量装载与数据校验 · 返回本章目录 · 下一节:异构同步的语义损失 · 查看全书目录 · 查看索引中心
29.5 异构同步的语义损失
异构同步可以让目标端“有数据”,却无法自动保证两边表达的是同一件事。connector 显示 running、offset 持续推进、目标查询也返回 200,只说明管道在工作;类型舍入、 排序规则、事务边界和删除语义仍可能已经变化。
本节给出一套语义合同。它不仅适用于 PostgreSQL 到 MySQL、Kafka、Elasticsearch 或 数据仓库,也适用于两个配置、扩展与 locale 不同的 PostgreSQL 环境。
29.5.1 类型、精度、排序规则与时区
类型映射必须是一张可测试的合同
不能只写:
至少要写清:
| 源语义 | 目标映射需要回答 |
|---|---|
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,而不是只测正常样本:
迁移前后都用同一个 canonical encoder 输出,比较规范化值和预期错误类别。若业务决定 允许损失,例如金额从 4 位小数舍入到 2 位,必须记录舍入规则、受影响行数、总误差和 批准人;不能让驱动默认转换替团队做决定。
时区问题常被样本掩盖
PostgreSQL 的 timestamptz 保存一个绝对时间点,显示受 session TimeZone 影响;
timestamp 不含时区。把前者格式化为本地字符串再写进后者,会永久丢掉 offset。
合同应明确:
还要验证 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 约束、事务顺序与删除语义
源端约束不会自动变成下游约束
源端可以依赖:
消息流通常只携带行变化,不携带这些证明。目标是搜索索引或对象存储时,甚至没有对应的 约束机制。于是“source 每次提交都合法”不能推出“sink 任意时刻都合法”。
例如源事务先创建 customer 再创建 order。若 connector 按 table 分 topic,下游并行 消费,order 可能先可见。解决方式不是祈祷消费者够快,而是明确:
- 是否保留 source transaction ID 和 commit boundary;
- 跨表事件是否要求原子可见;
- 不要求原子时,查询层如何隐藏未完成 batch;
- parent 缺失是重试、暂存、告警还是丢弃;
- checkpoint 在整个事务之后还是每条事件之后推进。
事务内 row order 也不能随意打散。账户扣款、入账和 ledger 三条事件若被三个 worker 独立提交,中间态会破坏守恒。高吞吐设计必须说明它牺牲了什么可见性,以及如何恢复。
upsert 需要版本,delete 需要墓碑
一个简单的:
只保证当前语句不因 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 可能比实时事件更晚到:
目标就回到了旧状态。每个写入路径都必须服从同一条条件:
重建期间可以使用新的目标 namespace/index/table,完成校验后原子交换;不要让不带 version 的历史 backfill 与实时流争写同一记录。
29.5.3 目标端可查询不等于语义等价
绿灯只能证明它声明的那一层
| 绿灯 | 能证明 | 不能证明 |
|---|---|---|
| connector running | 进程存活并执行主循环 | 没有跳过 poison event |
| offset advancing | 一些事件被确认 | sink 副作用完整、顺序正确 |
| target row count 相等 | 总行数一致 | 行内容、关联和删除一致 |
| target query 成功 | 语法和服务可用 | 排序、精度、完整性等价 |
| lag 接近零 | 消费接近 source head | 历史基线正确 |
异构验收应沿一条更强的梯子:
代表性查询不是随机挑十条 SELECT *,而应从业务清单中覆盖:
- equality、range、prefix、全文与排序;
- NULL、缺失字段、数组/JSON 嵌套;
- pagination 和 tie-breaker;
- 聚合、去重、金额与时区窗口;
- 删除、恢复、乱序和重复事件;
- 最大对象、热点 key 与大事务;
- 权限过滤和租户隔离。
每条都定义允许差异。例如搜索结果可能允许排名小幅变化,但不能跨租户;报表金额必须 精确相同;分析仓库允许 10 分钟最终一致,但 reconciliation 不允许永久缺口。
建立“允许损失登记表”
异构系统很少完全同构,现实做法不是假装零损失,而是让损失显式:
没有登记的差异一律视为 defect;登记项也要有 owner、验证和复审条件。这个机制防止 “已知差异”在口头交接中无限扩张。
权威源和修复方向必须唯一
持续 reconciliation 发现不一致时,先回答:
切流前通常以源端为权威;切流后目标已承接新写,不能继续无条件“用源覆盖目标”。 权威边界随迁移状态改变,必须随状态机一同记录。
本章的目标端 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 先打印并保存非秘密身份:
并断言 source 与 target system_identifier 不同。database 名称相同不等于同一个
数据源,IP 不同也不保证不是同一 cluster 的两个实例。
凭据按职责和环境隔离
至少拆分:
replication 角色需要源端 LOGIN REPLICATION,initial copy 还需要 published table
的 SELECT;它不需要目标端 DDL。目标 runtime 不应能回写源。验证角色不应因为要比较
数据而获得修复权限。
连接合同还包括:
- TLS mode、CA 与 hostname verification;
- 精确 HBA source CIDR 和 role classification;
CONNECT、schemaUSAGE、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 生成并组织:
以及需要 operator 落实的 source disable 和 re-routing 步骤。使用时应:
- 固定 Pigsty/模板版本;
- 审阅生成的变量、端点和 SQL;
- 把生成物与迁移 ticket 绑定;
- 在 disposable database 完整演练;
- 由实际应用/网络 owner 实施写围栏与切流;
- 用 PostgreSQL catalog 和应用身份反向复核。
本章正式实验在本地开发环境以 pg-test 为源、pg-meta 为目标,各自创建一次性
database 和角色。这只是为了在有限实验环境里证明真正的双 cluster 边界。生产中不应
因为教程这样做,就把承担 Pigsty 管理面的 pg-meta 当成默认业务迁移目标。
29.6.2 观察槽、WAL、延迟与切换流量
先用原生视图定义信号
源端:
目标端:
再加上 pg_subscription_rel table state、PostgreSQL log 和业务 marker。不同版本的
视图列会变化,自动化应先断言 server major,而不是对未知列 SELECT * 后按位置解析。
关键关系是:
正式实验禁用 subscription 后确认 slot inactive,生成 3,000 条变化:
重新启用后才追平。这个小规模实验不能给生产容量一个固定阈值,却证明“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 可以统一观察:
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 还是配置中心。 切流证据至少包括:
监控图中的目标 TPS 上升只是旁证。最强证据是使用真实 runtime 凭据建立一个新连接, 验证目标系统身份并完成可回读 canary。
29.6.3 保留源环境直到退出观察窗口
保留不是继续双主写入
切流后的源环境应进入受保护状态:
这样它既能支持调查和条件回退,又不会继续产生与目标分叉的新事实。若业务需要 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。退出流程应:
- 记录 subscription、publication、slot、owner 和最终位置;
- 停止/确认 consumer 不再需要;
- 在两端都可达时正常删除 subscription;
- 在源端确认 main slot 与 table-sync slot 均不存在;
- 再删除 publication 和迁移临时角色/权限;
- 轮换生产凭据;
- 将 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 双集群实验:
实验不是“看见几行复制过去”就结束,而要主动制造 consumer 停滞、显式 apply conflict 和不会报错的 silent drift,再完成 sequence 校准、写围栏、模拟切流、条件回退与精确 清理。
29.7.1 完成全量加增量同步
先读边界,再运行脚本
本章只允许在已确认的 Pigsty 四节点开发沙箱执行:
完整合同在
lab-contract.md,机器可读的环境、对象、数量和验收条件在
requirements.json,迁移状态与允许动作在
migration-contract.json,拓扑图在
topology.mmd。
实验只创建以下固定名称对象,并用随机 run_id marker 证明所有权:
它明确禁止读取现有业务表、修改 Pigsty inventory、修改 Patroni/持久参数、改变真实 HAProxy/PgBouncer/DNS/VIP、终止无关连接和 force-drop。marker、对象名或连接范围有一项 不匹配,runner 都失败关闭。
先做纯静态合同检查:
完整实验会创建和删除两个一次性数据库,必须在已确认沙箱中指定一个新的私有证据目录:
若希望分步审阅:
capture 在任何写入前检查:
- 两端 service、cluster、PostgreSQL major、primary 身份;
- 两个不同的 system identifier;
- source
wal_level = logical; - 目标 database、role、slot 和 subscription 起点不存在;
- 第 19、23、25、28 章上游证据存在且环境边界一致;
- 实验源文件散列与随后执行的版本一致;
- 远端临时目录和证据目录满足私有权限。
创建 schema 与逻辑复制对象
夹具包含:
两表都有稳定主键,orders.customer_id 引用 customer;状态、金额和更新时间具有固定
业务约束。runner 在源端创建 publication 和 logical slot,在目标端创建相同 schema 与
subscription,然后等待 pg_subscription_rel 中两张表都从同步状态进入 r。
正式参考 run 证明:
system identifier 是本次沙箱证据,不是读者环境中的预期常量。验收的是“两端不同且分别 绑定已声明 cluster”,不是数字本身。
让初始快照与持续变更汇合
初始复制完成后,源端执行:
流程等待目标读取 marker,再比较两端精确行数、有序摘要、金额合计、状态分布和业务 不变量。参考 run 的结果为:
这同时证明了 initial copy 与增量可以汇合,以及验证是在声明的 marker 边界之后执行。 它不证明生产大表所需时间,也不覆盖迁移期间的 DDL;后者仍需独立编排。
连接失败也是 preflight 结果
开发中的第一次候选 run 在 fixture 写入前被 HBA 拒绝,因为临时角色未匹配 Pigsty 的 group-role 分类。流程没有临时放宽认证,而是:
- 停止实验;
- 按 marker 清理两个 database、五个角色、slot/subscription;
- 证明所有临时对象不存在;
- 用
INHERIT FALSE, SET FALSE的成员关系满足 HBA 分类; - 重新从新的 run 和空证据目录开始。
生产演练也应如此:preflight 失败说明合同不成立,不能在原 run 上一边改权限一边继续, 否则最终证据无法说明实际执行了哪套安全边界。
29.7.2 注入消费者停滞与数据差异
反例一:consumer 停了,风险留在 source
runner 精确禁用 pg36_shop_sub,确认源端 pg36_shop_slot inactive,然后在源夹具中
生成固定 3,000 行变化。参考结果:
禁用动作只匹配本 run 的 subscription;负载行数固定,不改变
max_slot_wal_keep_size,也不制造无限 WAL。数字取决于 tuple、full-page image、
checkpoint 和版本,教学结论是方向:
重新启用并追平后必须再次校验 manifest,不能因 worker 恢复 running 就进入下一阶段。
反例二:显式冲突会停 apply
实验先在目标端插入:
再在源端提交同一 key、不同值。目标 apply 命中唯一键,PostgreSQL 18 的
subscription 统计出现 insert_exists:
此时 slot 仍可能存在,subscription 也仍是一个 catalog 对象,但 apply 已无法越过 冲突事务。修复流程必须先证明:
实验删除精确的目标冲突 fixture 行,让源事务重放,再等待 marker 与 manifest 收敛。 生产环境不能把“删除目标所有冲突行后重试”写成通用脚本;不同冲突可能代表合法的 target-only write。
反例三:静默漂移不会停 apply
runner 在目标端直接修改 order_id = 1。该行之后没有新的源变化,因此:
但 16 个稳定 hash bucket 中,只有 bucket 1 的行数/摘要不一致。实验处于切流前, 合同指定 source 为权威,因而按主键读取源行、记录修复前后摘要、修复目标并重跑所有 分桶,结果:
这组反例区分了两种故障:
| 故障 | apply 是否报错 | 主要发现方式 |
|---|---|---|
| 唯一键/缺行等显式 conflict | 通常会 | worker/log/pg_stat_subscription_stats |
| 目标手工写、错误 backfill 等 silent drift | 不一定 | 持续 manifest、分桶和业务不变量 |
运行状态与数据等价必须分别验收。
29.7.3 验证、切流、回退并输出迁移证据包
切流前修正 sequence 并建立写围栏
逻辑复制已经把 order_id = 900000 复制到目标,但 sequence 本身没有随 DML 推进:
runner 按源端最大 ID 校准目标 sequence。随后目标 runtime canary 得到 900001,证明
不会立刻与已迁移主键碰撞。
源端则撤销 runtime 的 DML 能力,用同一凭据实际发起 INSERT:
“执行过 revoke”不是证据,旧凭据的负向操作才是。生产还要盘点 owner、
SECURITY DEFINER、scheduler 和其他 writer。
只模拟路由,不碰真实平台
私有 route-history.json 记录:
它只是一台状态机的模拟输入。runner 会检查:
在模拟 target 阶段写入一条可识别 canary。回退前把这 1 条目标独占数据显式对账回源, 再切回 source 并写 rollback canary。最终参考结果:
这里证明的是“在目标独占写可枚举时,条件回退协议可执行”。真实业务流量已经写入目标后, 是否回退仍取决于 reverse sync、对账能力和不可补偿副作用。
证据包必须能反驳伪成功
私有证据目录包含:
review 会检查证据权限、schema、交叉字段、源文件 hash、私密信息与 public/private 边界。validator 不只验证成功样本,还要求:
也就是说,篡改 system identifier、初始行数、marker、retained WAL、conflict counter、
bucket repair、sequence、写围栏、路由边界或清理结论,都不能继续得到 pass。公开参考
摘要在
migration-run.json;它不含密码、conninfo、主机密钥
或原始行。
完成实验后可以对同一私有证据包重复审计:
但不能把另一个 run 的证据目录与当前源文件拼接使用。
清理也是验收阶段
正常删除目标 subscription 时,PostgreSQL 同时删除远端 main slot。runner 随后分别在 两端确认:
任何一项不成立,实验都不算完成。尤其不能为了让 CI 变绿而终止所有连接或删除所有 inactive slot。
从沙箱证据到生产迁移票据
本实验的最终决策仍是:
正式票据至少还要补:
- 真实 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、观察窗口和退出条件。
读者完成本节后,应能提交的不是一句“逻辑复制已同步”,而是一份可以回答同步了什么、 在哪个边界相等、故障怎样暴露、切流由谁执行、何时还能回退、如何证明已清理的迁移 证据包。
上一节:多集群迁移环境 · 返回本章目录 · 下一章:推陈出新:版本升级与回滚策略 · 查看全书目录 · 查看索引中心