如何使用psycopg2分批从PostgreSQL加载数据并保存为Parquet
PostgreSQL分批导出数据为Parquet格式(psycopg2实现)
你的核心问题是没有正确实现持续迭代获取数据批次的逻辑,当前代码只执行了一次fetchmany,且后续的循环写法错误(records是单次获取的结果集,不是迭代器,循环会遍历每条记录而非批次)。以下是修正后的完整实现方案:
问题分析
- 原代码仅调用一次
cursor.fetchmany(size=5),只会获取第一批5条数据,不会自动继续获取后续批次。 - 你期望的迭代逻辑中,
for batch in records会遍历records里的每一条记录,而非批次,导致每次循环处理的是单条数据而非整个批次。
正确实现代码
import psycopg2 import pandas as pd # 数据库连接参数 dbname = "your_db" user = "your_user" password = "your_pwd" host = "your_host" port = "your_port" query = "SELECT * FROM your_table;" batch_size = 5 batch_counter = 1 conn = psycopg2.connect(dbname=dbname, user=user, password=password, host=host, port=port) cursor = conn.cursor() cursor.execute(query) columns = [column[0] for column in cursor.description] # 循环获取批次数据,直到返回空列表 while True: records = cursor.fetchmany(size=batch_size) if not records: break # 没有更多数据,退出循环 # 将批次数据转为DataFrame并保存为Parquet df = pd.DataFrame(data=records, columns=columns) df.to_parquet(f"batch_{batch_counter}.parquet") print(f"已导出第 {batch_counter} 批次,共 {len(records)} 条数据") batch_counter += 1 # 关闭资源 cursor.close() conn.close()
关键说明
- 循环逻辑:通过
while True持续调用fetchmany,当返回空列表时说明所有数据已读取完毕,退出循环。 - 批次命名:用
batch_counter给每个Parquet文件加序号,避免文件名重复覆盖。 - 资源释放:最后务必关闭游标和数据库连接,避免资源泄漏。
额外优化建议
- 如果处理超大表,建议在查询中添加
ORDER BY子句,确保数据顺序稳定(避免不同批次数据乱序)。 - 可以根据实际需求调整
batch_size,过大可能导致内存占用过高,过小则会增加IO次数。
内容的提问来源于stack exchange,提问作者Любовь Пономарева
相关产品推荐
相关产品推荐

