逻辑解码支持DDL


难度 中等

逻辑解码需要源端数据库把DDL转换成某一种格式存储,可以是某种数据库语法格式的sql或者其他格式的文本文件,比如json。

而对于目标端数据库,不管是我们的数据库还是其他的数据库,我们都需要提供一种能把源端存储格式转换为标准sql的方式。

这个时候就需要将DDL转换为目标数据库的sql语法格式。

postgresql

postgresql逻辑解码原生不支持DDL,其逻辑复制流程如下:

WAL → logical decoding → output plugin → subscriber

WAL记录的是存储操作,并不是原始的SQL。

heap_insert
heap_update
heap_delete
INSERT INTO t ...
ALTER TABLE ...

逻辑解码的时候只会解析如下内容:

Relation OID
tuple data

最后通过系统 catalog 解析表结构。

在postgresql中,DDL 并不会进入 logical replication stream。postgreSQL 的逻辑复制要求:publisher 和 subscriber 必须 schema 完全一致。

postgresql的主要实现方式是pglogical。

gaussdb

alt text

polardb

polardb-for-postgresql扩展了pg的逻辑复制,使DDL能进入到logical replication。并通过 pubddl 参数控制复制范围。

CREATE PUBLICATION pub1
FOR ALL TABLES
WITH (pubddl='all');

todo:在源码中无法搜索到pubddl。

支持如下操作:

CREATE
ALTER
DROP
TRUNCATE

todo:是否还支持其他操作。

polardb实现流程如下:

DDL
 ↓
ProcessUtility 捕获
 ↓
LogLogicalMessage 写 WAL
 ↓
logical decoding decode message
 ↓
pgoutput 输出 message
 ↓
subscriber apply worker 执行 SQL

源码模块分布:

src/backend/tcop/
    utility.c                 ← 捕获DDL

src/backend/polar/
    polar_logical_ddl.c      ← DDL记录核心实现

src/backend/replication/logical/
    decode.c                 ← WAL decode
    reorderbuffer.c          ← change buffer

src/backend/replication/pgoutput/
    pgoutput.c               ← 输出DDL

src/backend/replication/logical/
    worker.c                 ← subscriber执行DDL
                PRIMARY

CREATE TABLE
     │
     ▼
event trigger
     │
     ▼
LogLogicalMessage
     │
     ▼
        WAL
   XLOG_LOGICAL_MESSAGE
     │
     ▼
logical decoding
     │
     ▼
pgoutput plugin
     │
     ▼
logical replication stream
     │
     ▼
           SUBSCRIBER
     │
     ▼
apply worker
     │
     ▼
SPI_execute(DDL)
src/backend/commands/event_trigger.c
src/backend/replication/logical/logical.c
src/backend/replication/logical/decode.c
src/backend/replication/logical/reorderbuffer.c
src/backend/replication/pgoutput/pgoutput.c
src/backend/replication/logical/worker.c
实现阶段 文件位置 主要职责
DDL 捕获 src/backend/commands/event_trigger.c 捕获 DDL 命令文本
WAL 写入 src/backend/replication/logical/logical.c 写入 logical message
WAL Decode src/backend/replication/logical/decode.c 解析 WAL logical message
Buffer src/backend/replication/logical/reorderbuffer.c 缓冲消息
Output src/backend/replication/pgoutput/pgoutput.c 输出到 logical replication
Apply src/backend/replication/logical/worker.c Subscriber apply SQL
Policy src/include/catalog/pg_publication.h 控制 DDL replication

主要函数

decrib

ProcessUtilitySlow

EventTriggerDDLCommandEnd

pg_event_trigger_ddl_commands

pg_logical_emit_message_text(PG_FUNCTION_ARGS)

pg_logical_emit_message_bytea

LogLogicalMessage

XLogRecPtr
LogLogicalMessage(const char *prefix, const char *message, size_t size,
          bool transactional)
{
  xl_logical_message xlrec;

  /*
   * Force xid to be allocated if we're emitting a transactional message.
   */
  if (transactional)
  {
    Assert(IsTransactionState());
    GetCurrentTransactionId();
  }

  xlrec.dbId = MyDatabaseId;
  xlrec.transactional = transactional;
  /* trailing zero is critical; see logicalmsg_desc */
  xlrec.prefix_size = strlen(prefix) + 1;
  xlrec.message_size = size;

  XLogBeginInsert();
  XLogRegisterData((char *) &xlrec, SizeOfLogicalMessage);
  XLogRegisterData(unconstify(char *, prefix), xlrec.prefix_size);
  XLogRegisterData(unconstify(char *, message), size);

  /* allow origin filtering */
  XLogSetRecordFlags(XLOG_INCLUDE_ORIGIN);

  return XLogInsert(RM_LOGICALMSG_ID, XLOG_LOGICAL_MESSAGE);
}
PG_RMGR(RM_LOGICALMSG_ID, "LogicalMessage", logicalmsg_redo, logicalmsg_desc, logicalmsg_identify, NULL, NULL, NULL, logicalmsg_decode, NULL, NULL, NULL)

ddl的类型XLOG_LOGICAL_MESSAGE

LogicalDecodingProcessRecord

ReorderBufferQueueMessage

action类型 REORDER_BUFFER_CHANGE_MESSAGE

pgoutput_change

pgoutput_message

subscriber

apply_dispatch

LOGICAL_REP_MSG_MESSAGE M


文章作者: 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