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

无法生成单CSV,如何从Python(Pandas)导入900万+行数据到PostgreSQL?

处理900万行Pandas数据导入PostgreSQL的可行方法

方法1:Pandas to_sql 分块导入(推荐)

直接用Pandas内置的to_sql方法,通过设置chunksize参数分批次插入,无需导出CSV。结合SQLAlchemy创建数据库连接引擎,性能远高于逐行插入。

示例代码:

import pandas as pd
from sqlalchemy import create_engine

# 创建PostgreSQL连接引擎
engine = create_engine('postgresql://username:password@host:port/dbname')

# 假设清洗后的数据存储在df中
df.to_sql(
    name='target_table',  # 目标表名
    con=engine,
    if_exists='append',  # 可选值:replace/append/fail
    index=False,
    chunksize=50000  # 每批次插入行数,根据内存情况调整为1万-10万
)

提示:搭配psycopg2-binary可提升性能,创建引擎时可添加connect_args={'options': '-c statement_timeout=0'}避免超时。

方法2:分块导出CSV + PostgreSQL COPY 命令

COPY是PostgreSQL原生高速导入命令,比INSERT快一个数量级。将大数据框分块导出为小CSV后,逐个用COPY导入。

步骤1:分块导出CSV

chunk_size = 1000000  # 每100万行生成一个CSV
for i, chunk in enumerate(df.groupby(df.index // chunk_size)):
    chunk[1].to_csv(f'data_chunk_{i}.csv', index=False, header=False)

步骤2:用psycopg2执行COPY命令

import psycopg2

conn = psycopg2.connect(
    dbname='dbname',
    user='username',
    password='password',
    host='host'
)
cur = conn.cursor()

# 逐个导入分块CSV
for i in range(len(df) // chunk_size + 1):
    with open(f'data_chunk_{i}.csv', 'r') as f:
        cur.copy_expert(
            "COPY target_table FROM STDIN WITH (FORMAT csv, DELIMITER ',', HEADER false)",
            f
        )
    conn.commit()

cur.close()
conn.close()

提示:若目标表未创建,可先执行df.head(0).to_sql('target_table', con=engine, index=False)生成空表结构。

方法3:psycopg2 execute_batch 批量插入

将数据分块转换为元组列表,用psycopg2.extras.execute_batch批量执行INSERT语句,性能优于普通executemany。

示例代码:

import psycopg2
from psycopg2.extras import execute_batch

conn = psycopg2.connect(
    dbname='dbname',
    user='username',
    password='password',
    host='host'
)
cur = conn.cursor()

# 编写插入语句(需与表列名对应)
insert_query = "INSERT INTO target_table (col1, col2, col3) VALUES (%s, %s, %s)"

# 分块处理数据
chunk_size = 50000
for i in range(0, len(df), chunk_size):
    chunk = df.iloc[i:i+chunk_size]
    # 将DataFrame转换为元组列表
    data_tuples = list(chunk.itertuples(index=False, name=None))
    execute_batch(cur, insert_query, data_tuples)
    conn.commit()

cur.close()
conn.close()

方法4:Dask处理超大数据(内存不足时)

若本地内存无法容纳900万行数据,可用Dask DataFrame替代Pandas,自动分块处理并批量导入。

示例代码:

import dask.dataframe as dd
from sqlalchemy import create_engine

# 将Pandas DataFrame转换为Dask DataFrame,分10个块处理
dask_df = dd.from_pandas(df, npartitions=10)

engine = create_engine('postgresql://username:password@host:port/dbname')

# 批量导入
dask_df.to_sql(
    name='target_table',
    uri=engine.url,
    if_exists='append',
    index=False,
    chunksize=50000
)

内容的提问来源于stack exchange,提问作者Evan W

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 12:43:28