You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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)
    

额外注意事项:

  1. 源库的发布PUB需要提前创建,且关联了你需要监听的表,复制用户需要有对应表的读取权限
  2. 你之前报错中的PUB,后缀多了逗号的问题,是低版本psycopg2的已知bug,升级psycopg2到2.9以上版本即可解决

内容的提问来源于stack exchange,提问作者Rashid Khalikov

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.24 15:06:04