逻辑复制解惑


难度 中等

在初次写入数据的时候,ddl写到了系统表,dml先写入日志,然后apply成数据。
当前逻辑复制的原理是根据lsn读取wal日志,将wal日志decode成sql,然后发送目标端。

因而类推到ddl,也是同样的,需要先读取ddl,将ddl记录decode成sql,然后发送给目标端。

  1. ddl是怎么保存到系统对象里的,这个有没有什么读取方式或者参考资料?
  2. 是否有一种方式可以直接从系统对象里解析出ddl?
  3. ddl应该是保存到数据库的元数据里,是不是有配套的读取写入接口?
  4. 直接从系统对象里解析ddl,然后decode发送,是不是理论上就能实现ddl同步?(表字段更新,定期同步最新表结构?)
  5. 从上面的逻辑看,ddl似乎跟事务没有关系?ddl同步跟事务的主要联系是什么?
  6. ddl按理说也会生成wal日志,也有对应的lsn,是否在wal日志上加一些标记或者记住对应表的wal日志的lsn就可以还原成sql?
  7. ddl和dml的依赖关系是根据事务号来的还是根据lsn来的?比如先执行ddl还是先执行dml的时间顺序是根据事务号还是lsn?

需要搞清楚的问题:

  1. 建立好复制关系之后,你分别做insert delete update,trunate,整个流程是什么样的?
  2. 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 中事务变更特点

  1. 每条 WAL 记录都是操作级别
  • INSERT / UPDATE / DELETE / TRUNCATE 等操作,WAL 记录中只包含:

    • 操作类型(XLOG_HEAP_INSERT / XLOG_HEAP_UPDATE 等)
    • relation OID(表标识)
    • tuple 数据/ctid(物理位置)
  • WAL 本身并不明确区分事务,但它包含 xid (TransactionId)

  1. 事务提交标记
  • 每个事务的 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 流程
  1. XLogSendLogical() 循环读取 WAL record。

  2. 每条 WAL record 调用 LogicalDecodingProcessRecord(ctx, record)

  3. 解析 WAL record 时:

    • 读取记录中的 xid(事务号)

    • 查询 ReorderBuffer 是否已经存在该事务的 entry:

      • 如果没有,创建一个新的 ReorderBufferTXN
      • 如果有,把本条变更加入该事务的 changes 列表
  4. 重复这个过程,直到遇到 XLOG_XACT_COMMIT

3. 提交事务时发送
  • 当 WAL record 是 commit:

    1. 找到该事务的 ReorderBufferTXN

    2. 调用 output plugin 回调:

      begin_cb(txn)
      change_cb(change1)
      change_cb(change2)
      ...
      commit_cb(txn)
    3. 释放 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

逻辑解码流程:

  1. 读取 100 → xid 10 → 放入 ReorderBuffer[10]
  2. 读取 110 → xid 11 → 放入 ReorderBuffer[11]
  3. 读取 120 → xid 10 → 添加到 ReorderBuffer[10]
  4. 读取 130 → xid 12 → 新建 ReorderBuffer[12]
  5. 读取 140 → xid 11 → 添加到 ReorderBuffer[11]
  6. 读取 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 等。
  • 解码流程:

    1. 解析 WAL,得到操作和表 OID。
    2. 查询 RelationSyncEntry,看这个表是否在 publication 中。
    3. 如果在 publication 中,就编码成逻辑复制消息(INSERT/UPDATE/DELETE)。
    4. 发送给 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

发布端在创建订阅时会:

  1. 创建 logical replication slot
  2. 获取一个 一致性快照

内部调用:

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_lsn
  • confirmed_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
=
最终一致数据

文章作者: growdu
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 growdu !
  目录
分类导航
随笔2 AI27 算法1 计算机基础13 博客搭建7 ChatGPT2 集群63 计算机通信1 数据库34 数据库深入80 DPDK26 Docker11 Elasticsearch4 编辑工具4 FAQ1 Go Web1 hometown2 编程语言16 网络9 OPC1 Linux38 openGauss4 页面12 PostgreSQL54 程序员自我修养1 协议11 成长之路1 stock1 存储5 工具20 VPP18 视频作品1 Vue13 Web1 代码示例11 数据库15 BenchmarkSQL1 PostgreSQL 源码修炼之路14
最热文章
1
13 逻辑复制深入
数据库深入🔥 1570
2
0 Postgresql存储、索引及系统优化、主备切换
PostgreSQL🔥 1495
3
一文读懂openguass dcf网络模块
集群🔥 1420
4
逻辑复制源码分析
数据库深入🔥 1327
5
PostgreSQL 分区表:从一行 `PARTITION BY` 到路由热路径的全链路拆解
数据库🔥 1094
6
applyparallelworker.c 之 LA 端源码深度解析:Leader Apply Worker 的指挥中枢
数据库深入🔥 1082
7
PostgreSQL Background Worker 全解:从 `RegisterBackgroundWorker` 到逻辑复制 4 类 worker 的全生命周期
数据库🔥 1078
8
PostgreSQL的后台进程walsender分析 - 关系型数据库 - 亿速云
PostgreSQL🔥 1033
9
PostgreSQL 逻辑复制的监控:六张视图 + 一组可执行 SQL,把 publisher/subscriber 的速率与健康度彻底看透
数据库🔥 1032
10
PostgreSQL 逻辑复制支持 DDL 之后:DDL 与 DML 的时序难题(重点:分区表)
数据库🔥 999
11
reorderbuffer.c 源码深度解析:PostgreSQL 逻辑复制的"事务重组引擎
数据库深入🔥 953
12
PostgreSQL 内核开发:读取一张表的 9 步标准流程与缓存全景
数据库🔥 938
13
从 `postgres` 二进制到生产级守护 —— PostgreSQL 最外层模块与启动全流程拆解
数据库🔥 936
14
支持逻辑复制同步 DDL 适配 SQL Server 方案
数据库深入🔥 934
15
PostgreSQL 逻辑复制的 ReorderBuffer 与事务机制:从一行 WAL 到一致性变更流的全链路绑定
数据库🔥 913
16
DDL同步架构(美化版)
数据库深入🔥 908
17
PostgreSQL Latch 机制详解:从一行 SetLatch 到 epoll 的内核之旅
数据库🔥 871
18
pgbench 源码全解:一个 C 文件如何撑起 PostgreSQL 官方压测工具
数据库🔥 860
19
PostgreSQL libpq 机制与缓冲区详解
数据库🔥 850
20
PostgreSQL 逻辑复制 spill 文件深度剖析:从 `xid-*.spill` 到 TPC-C 的增长方程
数据库🔥 845