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

使用pg8000驱动在Airflow DAG中执行Postgres COPY流时遇参数数据类型无法识别错误的求助

解决pg8000 + SQLAlchemy执行COPY FROM STDIN时的参数类型错误问题

你遇到的问题根源在于SQLAlchemy的conn.execute()方法并没有把stream参数传递给pg8000的底层cursor,而是将它当作普通的查询参数处理,这就导致PostgreSQL无法识别这个参数的数据类型,从而抛出could not determine data type of parameter $1的错误。

pg8000的原生cursor支持通过stream参数来处理COPY FROM STDIN操作,但SQLAlchemy的上层封装并没有适配这个特殊参数,所以我们需要直接使用pg8000的原生连接和cursor来执行COPY语句。

修改后的代码示例

def getconn() -> pg8000.native.Connection:
    conn: pg8000.native.Connection = connector.connect(
        PG_CONFIG["host"],
        "pg8000",
        user=PG_CONFIG["user"],
        password=PG_CONFIG["password"],
        db=PG_CONFIG["database"]
    )
    return conn

engine = sqlalchemy.create_engine("postgresql+pg8000://", creator=getconn)
engine.dialect.description_encoding = None

stream_in = StringIO()
csv_writer = csv.writer(stream_in)
csv_writer.writerow([1, "electron"])
csv_writer.writerow([2, "muon"])
csv_writer.writerow([3, "tau"])
stream_in.seek(0)

# 获取SQLAlchemy连接
with engine.connect() as conn:
    # 创建表的操作可以继续用SQLAlchemy的execute
    conn.execute("CREATE TABLE IF NOT EXISTS temp_table (user_id numeric, user_name text)")
    
    # 获取pg8000的原生连接和cursor
    native_conn = conn.connection
    cursor = native_conn.cursor()
    
    try:
        # 直接用原生cursor执行COPY语句,传递stream参数
        cursor.execute("COPY temp_table FROM STDIN WITH (FORMAT CSV)", stream=stream_in)
        # 提交事务(pg8000默认可能不是自动提交)
        native_conn.commit()
    finally:
        cursor.close()

关键说明

  • 通过conn.connection可以从SQLAlchemy的连接对象中获取到pg8000的原生连接实例。
  • 原生pg8000的cursor.execute()方法支持stream参数,它会正确处理COPY FROM STDIN的数据流,而不是把stream当作普通查询参数。
  • 记得手动提交事务,因为pg8000的连接默认不会自动提交修改,否则COPY的数据不会被持久化到数据库中。

你可以尝试这个修改方案,应该能解决你遇到的参数类型错误问题。

内容的提问来源于stack exchange,提问作者David Bayon

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 17:59:05