逻辑复制支持DDLAI提示词


难度 中等

需求

逻辑复制需要支持自动同步DDL。

逻辑复制当前仅支持同步DML,为了提升用户使用体验和增强逻辑同步功能。

publication

原来的publication用户接口:

CREATE PUBLICATION name
    [ FOR ALL TABLES
      | FOR publication_object [, ... ] ]
    [ WITH ( publication_parameter [= value] [, ... ] ) ]

where publication_object is one of:

    TABLE table_and_columns [, ... ]
    TABLES IN SCHEMA { schema_name | CURRENT_SCHEMA } [, ... ]

and table_and_columns is:

    [ ONLY ] table_name [ * ] [ ( column_name [, ... ] ) ] [ WHERE ( expression ) ]

需要增加ddl选项,接口变更为:

CREATE PUBLICATION name
    [ FOR ALL TABLES
      | FOR publication_object [, ... ] ]
    [ WITH ( publication_parameter [= value] [, ... ], ddl [ = value] [, ... ]) ]

where publication_object is one of:

    TABLE table_and_columns [, ... ]
    TABLES IN SCHEMA { schema_name | CURRENT_SCHEMA } [, ... ]

and table_and_columns is:

    [ ONLY ] table_name [ * ] [ ( column_name [, ... ] ) ] [ WHERE ( expression ) ]

其中,ddl可取的值如下:

  • table
  • index
  • type
  • function
  • domain
  • trigger
  • view
  • rule
  • schema
  • extension
  • all

其中table 和index在for table,for table in schema,for all tables下都生效,而function、domain、trigger、view、rule、schema、extension、all只在for all tables下生效,需要在实现时进行限制。

如果publication是FOR TABLE或FOR TABLES IN SCHEMA,ddl选项只能包含table和index。如果用户指定了function等类型,应该报错。

manual则是需要用户手动执行命令同步ddl,不会自动同步。应该提供一个函数如pg_sync_ddl(ddl)来手动触发同步。

subscription

原来的接口为:

CREATE SUBSCRIPTION subscription_name
    CONNECTION 'conninfo'
    PUBLICATION publication_name [, ...]
    [ WITH ( subscription_parameter [= value] [, ... ] ) ]

同样的增加ddl选项,接口变更为:

CREATE SUBSCRIPTION subscription_name
    CONNECTION 'conninfo'
    PUBLICATION publication_name [, ...]
    [ WITH ( subscription_parameter [= value] [, ... ], ddl [ = value] [, ... ] ) ]

ddl的取值范围与publication一致。publication用于控制发布,subscription用于控制是否订阅。

subscription在订阅时需要确定publication是否存在对应的发布,不存在需要报错。

设计实现

整体设计原则:

采用记录客户端ddl的query_string到新建的系统表pg_publication_sync,逻辑同步时如果开启了ddl同步,就把pg_publication_sync当作普通表使用dml同步过去。在该表中会记录ddl的原始字符串,订阅端收到后可以重新执行。

pg_publication_sync表的定义如下:

字段 类型 描述
lsn lsn 写入这个记录时,lsn位置,用于定时删除不需要的记录(做个函数清理pg_publication_sync表里不需要的信息: select pg_publication_sync_prune() ,由用户手动执行)
timestamp timestamp with time zone 事件上发生的时间
message_type char 消息类型,’A’ 表示新增对象、’D’表示删除对象、’Q’ DDL SQL信息,其它类型按需要扩展
message_data text 事件详细信息,根据message_type按需要定义即可
namespace text 当前search_path
publication text[] publication列表,需要根据这个列表过滤数据
ddl_type text[] 记录ddl的类型,订阅端需要根据这个字段对比本端的ddl类型,确认是否需要apply
target_table text 该ddl对应的table,如果是非table类的ddl,为null,主要是根据target table来确定该行记录是否需要发布
message_extra json 事件的额外信息,json格式,必选字段:版本、语法模式。其它可选参数按需添加:如一些特殊参数、大小写敏感等

提取ddl

可以完全参照log_statement=dll的实现(详见:check_log_statement),parseTree里会记录stmt对应的SQL文本,在query_string里的位置:standard_ProcessUtility的入参pstmt、queryString,得到DDL的原始SQL。

发布DDL

在standard_ProcessUtility中isCompleteQuery设置为true后发布DDL。即往pg_publication_sync中新增表。

  1. 发布ddl消息message_type=‘Q’
{
    "message_type": "Q",
    "namespace": "[mynsp,public,pg_catalog]",   --当前环境namespace
    "publication": "[pub1,pub2]",   --满足上面规则的publicaton
    "message_data": "CREATE TABLE t1(id int) ", --原始SQL
    "message_extra": "版本、语法模式、关键参数",   --apply时使用的必要的额外信息,json格式,按需添加
}

2.发步同步消息(message_type=‘A/D’)

{
    "message_type": "A",     --A表示新增,D表示删除(只有这种情况需要message_type='D'的消息)
    "publication": "[pub1,pub2]",  --被添加或者删除的pg_publication_rel.prpubid,且publication.ddl里有table类型
    "message_data": " {schema='nsp', table_name='tbl'} "  --新增的表对象,可以用json格式
}

pg_publication_sync需要跟随DML发布。

所有ddl都需要发布Q消息,A/D消息是用于发布变更,A消息细节可以参考:AlterSubscription_refresh,D消息可以参考AlterSubscription_refresh。

如何解码

发布端在发送数据时,需要解析wal日志,在读取到pg_publication_sync表数据时,需要把pg_publication_sync的行数据取出来,并进行过滤匹配。
以普通表的方式发送到订阅端。
需要注意的是,发布端在解码时会把系统表排除,需要根据pg_publication_sync的oid把这个表在解码流程中放行。

如何apply

订阅端在收到pg_publication_sync表的数据时需要进行特殊处理,识别到ddl消息后,将message_data提取出来执行。

订阅端直接复用DML的apply woker,然后按顺序解析执行就行,需要注意的是table sync worker不能处理DDL消息。

订阅端通过解析表名是pg_publication_sync来识别ddl同步,当发现是这张表时,就需要执行ddl同步逻辑。

主要工作事项

1.用户接口适配;
2.系统表实现;
3.发布端适配;
4.订阅端适配;
5.编译构建;
6.测试用例编写;
7.测试验证;
8.测试报告输出;


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