无法生成单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
相关产品推荐
相关产品推荐

