1、设置参数
gs_guc reload -D /database/panweidb/data/ -c 'wal_level=logical'
然后重启库生效
show wal_level2、确认用户有权限
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.log6、使用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进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




