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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 04:36:04