整体架构
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模式,整体流程如下:
- 后端进程捕获客户端的原始DDL存入系统表(相当于队列)
- 通过逻辑复制协议将DDL同步到订阅端
- 订阅端识别到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
× 不会跨方言捕获
分层架构
Capture
Rule Filter
Message Store
Replication Transport
复制传输层
Logical Decoding
pgoutput
Protocol
DDL Replay Framework
DDL 回放框架
Apply Worker
Rule Filter
Replay Engine
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
typedef enum ReplayDialect {
REPLAY_DIALECT_PG = 0,
REPLAY_DIALECT_TSQL,
REPLAY_DIALECT_ORACLE,
REPLAY_DIALECT_MYSQL,
REPLAY_DIALECT_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;
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,
};
static bool
pg_replay_execute_ddl(ReplayExecContext *ctx, const char *sql)
{
List *parsetree_list = pg_parse_query(sql);
return true;
}
SQL Server Replay
对于SQL Server来说,worker没有和后端进程的TDS端口建立连接,需要在worker进程启动时构建SQL Server parse上下文。构建上下文分为两种情况:
- 简单上下文:仅修改worker进程内的
sql_dialect,让捕获到的TSQL走bbf的parse
- 完整上下文:需要在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,
};
(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框架中。
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)
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)
{
return true;
}
static bool
mysql_replay_execute_ddl(ReplayExecContext *ctx, const char *sql)
{
return true;
}
void _PG_init(void)
{
RegisterReplayDialectAdapter(REPLAY_DIALECT_MYSQL, &mysql_adapter);
}
DDL Replication Platform
Replication Transport Layer
Capture Framework
Rule Filter Framework
Message Store
pg_publication_sync
L2: Replication Transport Layer
Logical Decoding
pgoutput
Replication Protocol
Apply Worker
Rule Filter
Replay Engine
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 到目标库