使用Python将S3数据导入RDS PostgreSQL时遇COPY文件不支持错误
将S3中2800个CSV文件导入RDS Postgres的问题与解决方法
问题背景
我尝试使用COPY命令将S3存储桶中的2800个CSV文件导入RDS Postgres,开发流程如下:
- 列出目标S3文件夹下的所有对象
- 在Postgres中创建对应结构的表
- 复制单个文件进行概念验证(POC)
原实现代码
import boto3 import psycopg2 S3_BUCKET = "arapbi" S3_FOLDER = "polygon/tickers/" s3 = boto3.resource("s3") my_bucket = s3.Bucket(S3_BUCKET) object_list = [] for obj in my_bucket.objects.filter(Prefix=S3_FOLDER): object_list.append(obj) conn_string = "postgresql://user:pass@db.address.us-west-2.rds.amazonaws.com:5432/arapbi" def write_sql(file): sql = f""" COPY tickers FROM '{file}' DELIMITER ',' CSV; """ return sql table_create_sql = """ CREATE TABLE IF NOT EXISTS public.tickers ( ticker varchar(20), timestamp timestamp, open double precision, close double precision, volume_weighted_average_price double precision, volume double precision, transactions double precision, date date )""" # 创建表 pg_conn = psycopg2.connect(conn_string, database="arapbi") cur = pg_conn.cursor() cur.execute(table_create_sql) pg_conn.commit() cur.close() pg_conn.close() # 尝试上传单个文件到表 sql_copy = write_sql(object_list[-1].key) pg_conn = psycopg2.connect(conn_string, database="arapbi") cur = pg_conn.cursor() cur.execute() pg_conn.commit() cur.close() pg_conn.close()
生成的COPY语句
COPY tickers FROM 'polygon/tickers/dt=2023-04-24/2023-04-24.csv' DELIMITER ',' CSV;
报错信息
运行COPY导入逻辑时,触发以下错误:
FeatureNotSupported: COPY from a file is not supported HINT: Anyone can COPY to stdout or from stdin. psql's \copy command also works for anyone.
修正后的解决方案代码
感谢Adrian Klaver和Thorn的指导,以下是更新后的可运行代码:
import pandas as pd from io import StringIO import psycopg2 conn_string = f"postgresql://{user}:{password}@arapbi20240406153310133100000001.c368i8aq0xtu.us-west-2.rds.amazonaws.com:5432/arapbi" pg_conn = psycopg2.connect(conn_string, database="arapbi") for i, file in enumerate(object_list): cur = pg_conn.cursor() output = StringIO() obj = object_list[i].key bucket_name = object_list[i].bucket_name df = pd.read_csv(f"s3a://{bucket_name}/{obj}").drop("Unnamed: 0", axis=1) output.write(df.to_csv(index=False, header=False, na_rep='NaN')) output.seek(0) cur.copy_expert(f"COPY tickers FROM STDIN WITH CSV HEADER", output) pg_conn.commit() cur.close() n_records = str(df.count()) print(f"已从s3a://{bucket_name}/{obj}加载{n_records}条记录") pg_conn.close()
注:原代码中enumerate(object_list[])存在语法错误,已修正为enumerate(object_list),同时补充了缺失的依赖导入语句。
内容的提问来源于stack exchange,提问作者Evan Volgas
相关产品推荐
相关产品推荐

