暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

panwei 之ddl+dml 记录回流

原创 feilunshuai 2026-06-26
106

1、设置参数

gs_guc  reload -D /database/panweidb/data/ -c 'wal_level=logical'
然后重启库生效
show wal_level


2、确认用户有权限

ALTER ROLE fei  WITH REPLICATION;
添加到pg_hba.conf中
gs_guc  reload -D /database/panweidb/data/ -h "host   replication     fei      192.168.0.0/24    sha256"
gs_guc  reload -D /database/panweidb/data/ -h "local   replication     fei                       trust"


3、创建逻辑复制槽

进到指定的数据库创建
SELECT * FROM pg_create_logical_replication_slot('my_slot', 'mppdb_decoding');


如果不需要删除
SELECT pg_drop_replication_slot('my_slot');


4、创建发布

CREATE PUBLICATION ydjf_pub FOR ALL TABLES with (ddl ='all'); 
如果是v3-3.4.0 以上版本可以换成
修改enable_ddl_logical_decode=on 


5、接收复制槽内容

pg_recvlogical -d testdb -U fei   --slot=my_slot  --start   --plugin=mppdb_decoding   --file=/tmp/pg_changes.log


6、使用python脚本解析为sql

#!/usr/bin/env python
# -*- coding:utf-8 -*-


from __future__ import print_function


import os
import json
import time
import re


LOG_FILE = "/tmp/pg_changes.log"
SQL_FILE = "/tmp/pg_changes.sql"
OFFSET_FILE = "/tmp/pg_changes.offset"


SLEEP = 0.2




########################################################################
# offset
########################################################################


def load_offset():


    if not os.path.exists(OFFSET_FILE):
        return 0


    try:
        with open(OFFSET_FILE) as f:
            return int(f.read().strip())
    except:
        return 0




def save_offset(offset):


    tmp = OFFSET_FILE + ".tmp"


    with open(tmp, "w") as f:
        f.write(str(offset))


    os.rename(tmp, OFFSET_FILE)




########################################################################
# value
########################################################################


def sql_value(v):


    if v is None:
        return "NULL"


    s = str(v)


    if s.lower() == "null":
        return "NULL"


    try:
        float(s)
        return s
    except:
        pass


    s = s.replace("'", "''")


    return "'%s'" % s




########################################################################
# DML
########################################################################


def parse_insert(obj):


    cols = obj["columns_name"]
    vals = obj["columns_val"]


    return "INSERT INTO %s(%s) VALUES (%s);" % (
        obj["table_name"],
        ",".join(cols),
        ",".join([sql_value(i) for i in vals])
    )




def parse_update(obj):


    sets = []


    for c, v in zip(obj["columns_name"], obj["columns_val"]):
        sets.append("%s=%s" % (c, sql_value(v)))


    where = []


    for c, v in zip(obj["old_keys_name"], obj["old_keys_val"]):
        where.append("%s=%s" % (c, sql_value(v)))


    if len(where) == 0:
        where.append("1=1")


    return "UPDATE %s SET %s WHERE %s;" % (
        obj["table_name"],
        ",".join(sets),
        " AND ".join(where)
    )




def parse_delete(obj):


    where = []


    for c, v in zip(obj["old_keys_name"], obj["old_keys_val"]):
        where.append("%s=%s" % (c, sql_value(v)))


    if len(where) == 0:
        where.append("1=1")


    return "DELETE FROM %s WHERE %s;" % (
        obj["table_name"],
        " AND ".join(where)
    )




########################################################################
# json
########################################################################


def parse_json(line):


    try:
        obj = json.loads(line)
    except:
        return None


    op = obj.get("op_type", "").upper()


    if op == "INSERT":
        return parse_insert(obj)


    elif op == "UPDATE":
        return parse_update(obj)


    elif op == "DELETE":
        return parse_delete(obj)


    return None




########################################################################
# ddl
########################################################################


DDL_RE = re.compile(r"decode to:\s*(.*?),\s*\[owner")




def parse_decode(line):


    m = DDL_RE.search(line)


    if not m:
        return None


    sql = m.group(1).strip()


    if not sql.endswith(";"):
        sql += ";"


    return sql




########################################################################
# output
########################################################################


def write_sql(sql):


    with open(SQL_FILE, "a") as f:
        f.write(sql)
        f.write("\n")




########################################################################
# main
########################################################################


def main():


    offset = load_offset()


    print("start offset =", offset)


    while True:


        if not os.path.exists(LOG_FILE):
            time.sleep(1)
            continue


        filesize = os.path.getsize(LOG_FILE)


        # logrotate
        if filesize < offset:
            offset = 0


        with open(LOG_FILE, "r") as f:


            f.seek(offset)


            while True:


                line = f.readline()


                if not line:
                    offset = f.tell()
                    save_offset(offset)
                    break


                offset = f.tell()


                line = line.strip()


                if line == "":
                    continue


                if line.startswith("BEGIN"):
                    continue


                if line.startswith("COMMIT"):
                    continue


                sql = None


                if line.startswith("{"):
                    sql = parse_json(line)


                elif "decode to:" in line:
                    sql = parse_decode(line)


                if sql:
                    write_sql(sql)


                save_offset(offset)


        time.sleep(SLEEP)




########################################################################


if __name__ == "__main__":
    main()


7、注意

该ddl 不包括truncate 操作。


「喜欢这篇文章,您的关注和赞赏是给作者最好的鼓励」
关注作者
【版权声明】本文为墨天轮用户原创内容,转载时必须标注文章的来源(墨天轮),文章链接,文章作者等基本信息,否则作者和墨天轮有权追究责任。如果您发现墨天轮中有涉嫌抄袭或者侵权的内容,欢迎发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论