如何将psycopg2逻辑复制与PostgreSQL的publisher发布者绑定
psycopg2逻辑复制绑定指定发布的实现方案
你遇到的报错核心原因是默认使用的test_decoding逻辑解码插件不支持publication_names参数,该参数是PostgreSQL内置的pgoutput解码插件(原生逻辑复制专用插件)的特有参数,按以下步骤操作即可实现绑定指定发布监听:
- 第一步:创建复制槽时指定使用
pgoutput插件
不要使用默认的test_decoding插件,创建逻辑复制槽的代码示例:import psycopg2 from psycopg2 import extras conn = psycopg2.connect( dbname="你的源库名", user="有replication权限的用户", password="密码", host="源库地址", port="源库端口", replication="database" # 必须指定为数据库复制模式 ) cur = conn.cursor() # 创建pgoutput类型的复制槽,若槽已存在可跳过这步 try: extras.create_logical_replication_slot(cur, 'pytest', 'pgoutput') except psycopg2.errors.DuplicateObject: cur.execute("ROLLBACK;") - 第二步:调用start_replication时正确传参
除了指定publication_names,还要注意decode参数不要设为True,pgoutput返回的是二进制格式的复制消息,需要按协议解析:cur.start_replication( slot_name='pytest', decode=False, # pgoutput模式下必须设为False options={ 'publication_names': 'PUB', 'proto_version': '1' # pgoutput必须指定协议版本,目前固定为1 } ) - 第三步:处理接收到的复制消息
pgoutput返回的消息需要按照PostgreSQL逻辑复制协议解析,示例代码如下:def consume_msg(msg): # 可自行解析msg.payload获取DML/DDL操作内容,也可引入第三方pgoutput解析库简化逻辑 print(f"收到变更: {msg.payload}") msg.cursor.send_feedback(flush_lsn=msg.data_start) cur.consume_stream(consume_msg)
额外注意事项:
- 源库的发布
PUB需要提前创建,且关联了你需要监听的表,复制用户需要有对应表的读取权限- 你之前报错中的
PUB,后缀多了逗号的问题,是低版本psycopg2的已知bug,升级psycopg2到2.9以上版本即可解决
内容的提问来源于stack exchange,提问作者Rashid Khalikov
相关产品推荐
相关产品推荐

