DDL同步架构(美化版)


难度 中等

整体架构

DDL Synchronization Architecture
Publisher Side
ProcessUtility()
CapturePublicationSyncDDL()
DDL Rule Filter
Type Schema Regex Plugin
pg_publication_sync
Logical Replication
Logical Decoding
pgoutput
Replication Protocol
Subscriber Side
apply_handle_ddl()
DDL Rule Filter
subddl Security Object Regex
Replay Framework
execute_publication_sync_sql_command()
PostgreSQL Adapter
Babelfish Adapter
Oracle Adapter
MySQL Adapter
核心特性
  • 统一架构:复用内核逻辑复制能力,最小化改造
  • 双向过滤:发布端+订阅端双重规则过滤
  • 多方言适配:按目标数据库方言原样执行DDL,不依赖语法转换
  • 基于数据表的过滤规则:支持按表、Schema、类型等多维度过滤

设计背景

当前流程:逻辑复制DDL同步已适配PG模式,整体流程如下:

  1. 后端进程捕获客户端的原始DDL存入系统表(相当于队列)
  2. 通过逻辑复制协议将DDL同步到订阅端
  3. 订阅端识别到DDL后,对DDL进行replay(使用原始DDL SQL走parse再execute执行)

核心挑战:replay最重要的就是apply worker需要有对应数据库模式的上下文。

  • 对于PG模式:sql解析引擎是PG,天生具有对应的DDL执行上下文,可直接调用pg_parse_query执行
  • 对于SQL Server模式:使用bbf插件实现,采用TCP端口区分,worker进程没有数据库连接,当前机制无法构造完整的SQL Server执行上下文
  • 考虑到后续可能适配MySQL、Oracle、DB2等模式,需要一种通用的replay框架方便扩展

关键约束

传输端口约束

  • 逻辑复制只能通过 PG 端口进行同步
  • 创建订阅时,指定的发布端端口必须是 PG 端口
  • 不能使用 TDS 端口或其他数据库方言端口作为复制传输通道
  • 原因:逻辑复制协议是 PostgreSQL 原生能力,依赖于 PG 端口的协议栈

DDL 捕获匹配规则

  • 只捕获与当前数据库模式匹配的 DDL 语句
  • 例如:SQL Server 模式下不会捕获 PG SQL
  • PG 模式下不会捕获 T-SQL
  • 原因:不同方言的 DDL 语法不同,捕获非匹配方言的语句会导致 replay 失败
PG 端口传输
创建订阅 → 连接 PG 端口 → 逻辑复制协议
× 不能使用 TDS 端口
方言匹配捕获
TSQL 模式 → 仅捕获 T-SQL DDL
PG 模式 → 仅捕获 PG DDL
× 不会跨方言捕获

分层架构

DDL Producer
DDL 生产者
Capture
Rule Filter
Message Store
Replication Transport
复制传输层
Logical Decoding
pgoutput
Protocol
DDL Replay Framework
DDL 回放框架
Apply Worker
Rule Filter
Replay Engine
Adapter Framework
适配器框架
Babelfish
Adapter
Oracle
Adapter
MySQL
Adapter
DB2
Adapter
分层架构优势
  • 职责分离:各层专注自身职责,便于维护和扩展
  • 标准接口:层间通过标准接口通信,解耦设计
  • 插件式适配器:新增数据库支持只需实现新适配器
  • 可测试性:各层可独立测试,保证质量

Replay Framework

Replay 架构总览

Apply Worker 收到 DDL 消息
Replay Framework 调度层
PG Replay
直接调用 pg_parse_query
TSQL Simple
设置 sql_dialect 后 parse
TSQL Custom
bbf插件注册完整上下文
Extension
MySQL/Oracle/DB2...
目标数据库执行 DDL

数据库方言定义

PostgreSQL
REPLAY_DIALECT_PG
T-SQL (SQL Server)
REPLAY_DIALECT_TSQL
Oracle
REPLAY_DIALECT_ORACLE
MySQL
REPLAY_DIALECT_MYSQL
DB2
REPLAY_DIALECT_DB2
typedef enum ReplayDialect { REPLAY_DIALECT_PG = 0, // PostgreSQL REPLAY_DIALECT_TSQL, // SQL Server (T-SQL) REPLAY_DIALECT_ORACLE, // Oracle REPLAY_DIALECT_MYSQL, // MySQL REPLAY_DIALECT_DB2, // IBM DB2 REPLAY_DIALECT_MAX } ReplayDialect;

Replay 上下文

ReplayExecContext 结构体

dialect
数据库方言类型 (ReplayDialect)
dbname
逻辑数据库名称
search_path
搜索路径设置
user_oid
用户对象ID
switched_user
是否已切换用户
ucxt
用户上下文 (UserContext)
adapter_private
扩展数据指针 (void*)
typedef struct ReplayExecContext { ReplayDialect dialect; // 数据库模式 char *dbname; // 逻辑数据库名称 char *search_path; // 搜索路径 Oid user_oid; // 用户对象ID bool switched_user; // 是否已切换用户 UserContext ucxt; // 用户上下文 void *adapter_private; // 扩展数据(各适配器私有) } ReplayExecContext;

接口定义

ReplayDialectAdapter 回调接口

每个数据库方言适配器必须实现以下三个核心回调函数:

name
适配器名称标识
init_context(ctx)
初始化replay执行上下文,返回bool表示成功/失败
execute_ddl(ctx, sql)
执行DDL语句,在已构建的上下文中parse并execute
cleanup_context(ctx)
清理replay执行上下文,释放资源
typedef struct ReplayDialectAdapter { const char *name; bool (*init_context)(ReplayExecContext *ctx); bool (*execute_ddl)(ReplayExecContext *ctx, const char *sql); void (*cleanup_context)(ReplayExecContext *ctx); } ReplayDialectAdapter;

Replay Framework 注册初始化接口

ReplayFrameworkInit()
框架全局初始化,注册内置适配器
RegisterReplayDialectAdapter(dialect, adapter)
注册/覆盖某个方言的适配器实现
GetReplayAdapter(dialect)
根据方言类型获取对应适配器
// 框架全局初始化 void ReplayFrameworkInit(void); // 注册方言适配器 bool RegisterReplayDialectAdapter(ReplayDialect dialect, ReplayDialectAdapter *adapter); // 获取方言适配器 ReplayDialectAdapter* GetReplayAdapter(ReplayDialect dialect);

接口使用流程

1

框架初始化

Postmaster启动时调用 ReplayFrameworkInit(),注册PG内置适配器和Simple TSQL适配器

2

插件注册(可选)

Babelfish插件在 _PG_init() 中调用 RegisterReplayDialectAdapter() 注册Custom TSQL适配器,覆盖Simple模式

3

DDL消息到达

Apply Worker收到DDL消息后,从消息中解析出 ReplayDialect 类型

4

获取适配器

调用 GetReplayAdapter(dialect) 获取对应方言的适配器实现

5

执行DDL

依次调用 init_context()execute_ddl()cleanup_context()

Replay 执行流程

Apply Worker 收到DDL消息
GetReplayAdapter(dialect)
adapter->init_context(ctx) 构建执行上下文(PG跳过,TSQL设置dialect/bbf上下文)
adapter->execute_ddl(ctx, sql) pg_parse_query → 执行计划 → ExecutorRun
adapter->cleanup_context(ctx) 恢复GUC设置、释放资源、清理adapter_private
DDL回放完成

PostgreSQL Replay

对于PG replay来说,只是封装一下replay ddl函数即可,不需要初始化和清理上下文。因为PG的apply worker天然具有PG执行上下文。

static ReplayDialectAdapter pg_adapter = { .name = "postgres", .init_context = pg_replay_init_context, .execute_ddl = pg_replay_execute_ddl, .cleanup_context = pg_replay_cleanup_context, }; // PG实现:直接调用pg_parse_query,无需额外上下文设置 static bool pg_replay_execute_ddl(ReplayExecContext *ctx, const char *sql) { List *parsetree_list = pg_parse_query(sql); // ... parse -> plan -> execute return true; }

SQL Server Replay

对于SQL Server来说,worker没有和后端进程的TDS端口建立连接,需要在worker进程启动时构建SQL Server parse上下文。构建上下文分为两种情况

  1. 简单上下文:仅修改worker进程内的 sql_dialect,让捕获到的TSQL走bbf的parse
  2. 完整上下文:需要在bbf内部设置逻辑数据库,使worker执行DDL与客户端连接TDS端口执行一致
对比项 Simple Replay Custom (Full) Replay
实现位置 内核源码 bbf插件源码
DDL支持范围 有限(基础DDL) 更完整(复杂DDL)
上下文构建 仅设置 sql_dialect 完整bbf上下文+逻辑数据库
是否需要改插件 不需要 需要实现并注册回调
注册方式 内核默认注册 插件 _PG_init() 中注册
执行一致性 部分一致 与TDS端口执行完全一致

Simple Replay 实现

static ReplayDialectAdapter tsql_simple_adapter = { .name = "tsql-simple", .init_context = tsql_simple_init_context, .execute_ddl = tsql_simple_execute_ddl, .cleanup_context = tsql_simple_cleanup_context, }; // Simple方式核心:设置sql_dialect让bbf的parser生效 (void) set_config_option("babelfishpg_tsql.sql_dialect", "tsql", PGC_USERSET, PGC_S_SESSION, GUC_ACTION_SET, true, 0, false);

Custom (Full) Replay 实现

完整的SQL Server上下文构建需要在bbf插件内部注册对应的回调函数,并在bbf插件初始化时注册到replay框架中。

// 在babelfishpg_tsql插件的_PG_init_函数中注册 void bbf_register_tsql_replay_adapter(void) { RegisterReplayDialectAdapter(REPLAY_DIALECT_TSQL, &bbf_tsql_adapter); get_current_dbname_hook = bbf_get_current_dbname; } static ReplayDialectAdapter bbf_tsql_adapter = { .name = "tsql", .init_context = bbf_replay_init_context, .execute_ddl = bbf_replay_execute_ddl, .cleanup_context = bbf_replay_cleanup_context, };

运行机制

Simple vs Custom 回退机制

1
Simple Replay默认注册:内核启动时自动注册Simple TSQL适配器
2
插件检查:如果bbf插件实现了对应的replay回调函数并注册到框架
3
覆盖注册:Custom适配器覆盖Simple适配器,不再使用Simple模式
4
降级保障:如果bbf插件未实现回调,自动降级为Simple Replay模式执行DDL
内核启动 → 注册 Simple TSQL Adapter
Babelfish插件加载 调用 _PG_init()
是否实现了 bbf_register_tsql_replay_adapter()?
Yes 注册 Custom Adapter 覆盖 Simple Adapter → 完整上下文执行
No 保持 Simple Adapter 使用默认实现 → 简单上下文执行

扩展新数据库模式

新增数据库模式(如MySQL、Oracle、DB2)只需实现以下步骤:

1

添加方言枚举

ReplayDialect 枚举中添加新类型(如 REPLAY_DIALECT_MYSQL

2

实现三个回调函数

init_context()execute_ddl()cleanup_context()

3

构造适配器结构体

填充 ReplayDialectAdapter 结构体,设置name和三个函数指针

4

注册到框架

在插件初始化时调用 RegisterReplayDialectAdapter(REPLAY_DIALECT_XXX, &xxx_adapter)

// 扩展示例:MySQL适配器 static ReplayDialectAdapter mysql_adapter = { .name = "mysql", .init_context = mysql_replay_init_context, .execute_ddl = mysql_replay_execute_ddl, .cleanup_context = mysql_replay_cleanup_context, }; static bool mysql_replay_init_context(ReplayExecContext *ctx) { // 设置MySQL dialect,初始化parser上下文 // 例如:设置兼容模式、字符集等 return true; } static bool mysql_replay_execute_ddl(ReplayExecContext *ctx, const char *sql) { // 调用MySQL兼容层的parse和execute // 例如:mysql_parse_query(sql) -> execute return true; } // 在插件初始化时注册 void _PG_init(void) { RegisterReplayDialectAdapter(REPLAY_DIALECT_MYSQL, &mysql_adapter); }

DDL Replication Platform

DDL Producer Layer
Replication Transport Layer
DDL Replay Framework
Adapter Framework
L1: DDL Producer Layer
Capture Framework
Rule Filter Framework
Message Store
pg_publication_sync
L2: Replication Transport Layer
Logical Decoding
pgoutput
Replication Protocol
L3: DDL Replay Framework
Apply Worker
Rule Filter
Replay Engine
L4: Adapter Framework
PostgreSQL
Babelfish
Oracle
MySQL
平台亮点
  • 复用内核能力:利用 PostgreSQL 原生逻辑复制能力
  • 高扩展:四层架构,每层可独立演进
  • 多支持:同时支持 Babelfish/Oracle/MySQL
  • 易集成:与现有复制体系无缝融合

Dual Rule Filter Framework

DDL SQL
Publisher
Rule Filter
DDL Message
Subscriber
Rule Filter
Replay Engine
双重过滤机制
  • 发布端过滤:控制哪些 DDL 需要同步,支持类型/Schema/正则/插件规则
  • 订阅端过滤:安全校验、对象存在性检查、冲突检测
  • 配置灵活:两端规则可独立配置,按需调整
  • 调试方便:双重过滤便于追踪 DDL 流转全过程

DDL Replication Platform

DDL SQL
Capture Framework
Rule Filter Framework
DDL Message
Replication Transport
Replay Framework
Adapter Framework
Target Database
端到端流程
  • 捕获:ProcessUtility 钩子捕获 DDL 语句
  • 过滤:规则引擎过滤需要同步的 DDL
  • 传输:逻辑复制协议将消息送达订阅端
  • 回放:适配器按方言原样执行 DDL 到目标库

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