在初次写入数据的时候,ddl写到了系统表,dml先写入日志,然后apply成数据。
当前逻辑复制的原理是根据lsn读取wal日志,将wal日志decode成sql,然后发送目标端。
因而类推到ddl,也是同样的,需要先读取ddl,将ddl记录decode成sql,然后发送给目标端。
- ddl是怎么保存到系统对象里的,这个有没有什么读取方式或者参考资料?
- 是否有一种方式可以直接从系统对象里解析出ddl?
- ddl应该是保存到数据库的元数据里,是不是有配套的读取写入接口?
- 直接从系统对象里解析ddl,然后decode发送,是不是理论上就能实现ddl同步?(表字段更新,定期同步最新表结构?)
- 从上面的逻辑看,ddl似乎跟事务没有关系?ddl同步跟事务的主要联系是什么?
- ddl按理说也会生成wal日志,也有对应的lsn,是否在wal日志上加一些标记或者记住对应表的wal日志的lsn就可以还原成sql?
- ddl和dml的依赖关系是根据事务号来的还是根据lsn来的?比如先执行ddl还是先执行dml的时间顺序是根据事务号还是lsn?
需要搞清楚的问题:
- 建立好复制关系之后,你分别做insert delete update,trunate,整个流程是什么样的?
- create subcription的时候有个copy data的选项,这时发布端如果有insert delete update,最后整个表怎么保证一致性
1️事务提交与 WAL 关系
- 每个事务在 commit 时都会在 WAL 中写入 XLOG_XACT_COMMIT 类型的记录。
- WAL 中的变更(INSERT/UPDATE/DELETE/Truncate)本身是按照事务执行顺序写入的。
- 逻辑解码并不是逐条发送 WAL,而是把 WAL 中同一事务的所有变更收集起来,然后在 commit 时统一发送给 subscriber。
- 所以,逻辑解码的 事务顺序 = WAL 中 commit LSN 顺序,保证 subscriber 接收顺序和事务提交顺序一致。
LSN 推进
- 逻辑解码会按照 LSN(WAL 位置) 依次解析 WAL 记录。
- 每解析一条 WAL,就知道对应的操作和表 OID。
- 对于交错事务,逻辑解码会把同一事务的操作缓存到 ReorderBuffer。
- 当 WAL 中出现 commit record 时,逻辑解码把整个事务(ReorderBuffer 中的 changes)一次性发出去。
如何将同一事务的操作合并
WAL 本身是物理层面的记录,没有事务概念,但逻辑复制通过 ReorderBuffer 把同一事务的 WAL 变更聚合起来。
WAL 中事务变更特点
- 每条 WAL 记录都是操作级别
INSERT / UPDATE / DELETE / TRUNCATE 等操作,WAL 记录中只包含:
- 操作类型(XLOG_HEAP_INSERT / XLOG_HEAP_UPDATE 等)
- relation OID(表标识)
- tuple 数据/ctid(物理位置)
WAL 本身并不明确区分事务,但它包含 xid (TransactionId)
- 事务提交标记
- 每个事务的 commit 在 WAL 中写入 XLOG_XACT_COMMIT
- abort 会写 XLOG_XACT_ABORT
所以 WAL 中虽然交错了多个事务的操作,但每条记录知道属于哪个事务号 (xid)
ReorderBuffer 收集同一事务变更
1. ReorderBuffer 结构
typedef struct ReorderBufferTXN
{
TransactionId xid;
List *changes; // 保存该事务的所有逻辑变更
XLogRecPtr commit_lsn; // commit LSN
bool committed;
} ReorderBufferTXN;
2. 逻辑解码解析 WAL 流程
XLogSendLogical()循环读取 WAL record。每条 WAL record 调用
LogicalDecodingProcessRecord(ctx, record)。解析 WAL record 时:
读取记录中的 xid(事务号)
查询 ReorderBuffer 是否已经存在该事务的 entry:
- 如果没有,创建一个新的
ReorderBufferTXN - 如果有,把本条变更加入该事务的
changes列表
- 如果没有,创建一个新的
重复这个过程,直到遇到 XLOG_XACT_COMMIT
3. 提交事务时发送
当 WAL record 是 commit:
找到该事务的
ReorderBufferTXN调用 output plugin 回调:
begin_cb(txn) change_cb(change1) change_cb(change2) ... commit_cb(txn)释放 ReorderBuffer 记录(如果未 spill to disk)
核心:逻辑解码并不是依赖 WAL 顺序发送,而是依赖 ReorderBuffer 聚合同一事务的所有 WAL 变更,直到 commit 才发送。
事务重组示例
假设 WAL 顺序如下(WAL 交错多个事务):
| WAL LSN | XID | 操作 |
|---|---|---|
| 100 | 10 | INSERT t1 |
| 110 | 11 | INSERT t2 |
| 120 | 10 | UPDATE t1 |
| 130 | 12 | DELETE t3 |
| 140 | 11 | UPDATE t2 |
| 150 | 10 | COMMIT |
逻辑解码流程:
- 读取 100 → xid 10 → 放入 ReorderBuffer[10]
- 读取 110 → xid 11 → 放入 ReorderBuffer[11]
- 读取 120 → xid 10 → 添加到 ReorderBuffer[10]
- 读取 130 → xid 12 → 新建 ReorderBuffer[12]
- 读取 140 → xid 11 → 添加到 ReorderBuffer[11]
- 读取 150 → xid 10 COMMIT → 输出 ReorderBuffer[10] 的所有变更
这样即使 WAL 交错写入,逻辑复制也能保证事务内部顺序和完整性。
- WAL 每条记录都有 xid,逻辑复制利用这个 xid 聚合事务变更
- ReorderBuffer 负责收集同一事务的所有 WAL 变更
- 事务 commit 时一次性触发 output plugin 发送
- 这样保证了 按事务提交顺序、事务内部操作顺序发送给 subscriber
- 对于大事务,ReorderBuffer 会 spill 到磁盘 仍保持完整性
表 OID 与 Publication 对比
WAL 记录本身包含 relation OID。
逻辑解码插件(如 pgoutput)维护 Publication + RelationSyncEntry:
- Publication 定义了哪些表的变更需要发送。
- RelationSyncEntry 缓存了每个表的列信息、OID 等。
解码流程:
- 解析 WAL,得到操作和表 OID。
- 查询 RelationSyncEntry,看这个表是否在 publication 中。
- 如果在 publication 中,就编码成逻辑复制消息(INSERT/UPDATE/DELETE)。
- 发送给 subscriber。
发送给 subscriber 的顺序
- 每个事务的变更按 事务内部顺序发送。
- 多个事务按 commit LSN 顺序发送。
- 这样 subscriber 可以按照逻辑顺序重放,保证一致性。
可以简化为:
事务提交 → WAL 写入 → 逻辑解码按 LSN 解析 WAL → 收集到 ReorderBuffer
→ commit 时触发 plugin 回调 → 检查表是否在 publication → 如果在则发送逻辑变更消息
所以:
- 顺序 = commit LSN 顺序
- 表过滤 = WAL relation OID + publication 对比
- 事务完整性 = ReorderBuffer + commit 时发送
- 断开重连一致性 = replication slot + confirmed_flush_lsn
copy_data = true,初始表拷贝 + 并发事务如何保证最终一致
通过 snapshot + replication slot LSN 边界 + 后续 WAL streaming 来保证一致性。
copy_data = true 的整体流程
创建订阅时的整体流程:
Subscriber
|
| CREATE SUBSCRIPTION
v
Publisher
|
| 创建 logical replication slot
| 获取 consistent snapshot
|
| - snapshot LSN -
|
| COPY table data
|
| 同时 WAL 继续产生
|
| COPY 完成
|
| 从 snapshot LSN 开始 streaming WAL
关键点:
COPY 时使用一个一致性 snapshot。
关键机制 1:Snapshot
发布端在创建订阅时会:
- 创建 logical replication slot
- 获取一个 一致性快照
内部调用:
SnapBuildExportSnapshot()
得到:
snapshot_lsn
这个 snapshot 表示:
“数据库在 snapshot_lsn 时刻的一致视图”
COPY 数据时发生的事情
COPY 不是简单的 SELECT * FROM table。
它是:
SET TRANSACTION SNAPSHOT snapshot_lsn
COPY table
因此:
- COPY 看到的数据 只包含 snapshot_lsn 之前提交的事务
- snapshot_lsn 之后提交的事务 不会被 COPY 看到
COPY 期间发生 INSERT/UPDATE/DELETE 怎么办?
假设流程:
T1: snapshot_lsn = 100
T2: 开始 COPY table
T3: 事务 A 在 LSN 110 INSERT
T4: 事务 B 在 LSN 120 DELETE
T5: COPY 结束
T6: 开始 WAL streaming
COPY 看到的数据:
<= LSN 100
事务 A/B:
> LSN 100
所以:
COPY 不包含 A/B
但是:
WAL streaming 会从 LSN 100 开始发送
因此:
subscriber 会收到:
INSERT (A)
DELETE (B)
最终结果:
COPY snapshot
+
WAL replay
=
一致数据
关键机制 2:Replication Slot
逻辑复制 slot 记录:
restart_lsnconfirmed_flush_lsn
创建 slot 时:
restart_lsn = snapshot_lsn
含义:
WAL 从 snapshot_lsn 开始 必须保留
否则 subscriber 可能追不上。
每个表的同步 worker
PostgreSQL 实际实现更复杂:
每个表有独立同步 worker:
table sync worker
流程:
1 创建临时 slot
2 COPY table
3 同步 snapshot
4 catch-up WAL
5 切换到主 slot
源码在:
src/backend/replication/logical/tablesync.c
核心函数:
LogicalRepSyncTableStart()
表同步阶段(更精确)
表同步有 4 个状态:
| 状态 | 含义 |
|---|---|
| INIT | 准备同步 |
| DATASYNC | COPY 数据 |
| CATCHUP | 追 WAL |
| SYNCDONE | 完成 |
关键阶段:
DATASYNC
COPY table
CATCHUP
同步 worker 会:
start replication from snapshot_lsn
追 WAL 直到:
current_lsn
这样保证:
COPY期间产生的所有变化
都被 replay
最终切换
当表追到最新 WAL:
tablesync worker 退出
主 apply worker 接管。
举一个完整例子
发布端:
table users
初始数据:
id
1
2
3
创建订阅:
snapshot_lsn = 100
COPY 开始:
COPY users -> 1 2 3
COPY 期间:
INSERT 4 LSN 110
DELETE 2 LSN 120
UPDATE 3 LSN 130
subscriber 收到:
第一步 COPY
1
2
3
第二步 WAL
INSERT 4
DELETE 2
UPDATE 3
最终:
1
3(updated)
4
与发布端一致。
为什么不会丢数据
因为:
snapshot_lsn
+
WAL replay
=
完整历史
只要满足:
WAL >= snapshot_lsn
就不会丢。
Replication slot 保证:
WAL 不会被删除
如果 COPY 很慢会怎样
如果 COPY 很慢:
WAL backlog 会增长
因为:
slot restart_lsn = snapshot_lsn
WAL 不会回收。
所以生产环境常见问题:
logical replication slot 导致 WAL 堆积
copy_data=true 的一致性机制:
一致性 snapshot
+
从 snapshot_lsn 开始的 WAL replay
=
最终一致数据