使用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
相关产品推荐
相关产品推荐

