使用Python Polars存储大DataFrame到PostgreSQL报错求解决方案
Polars写入PostgreSQL大DataFrame报错修复方案
错误根源
SQLSTATE: 22P04错误是PostgreSQL的COPY命令无法识别数据签名,本质是大文件传输时ADBC驱动的批量写入缓冲区溢出,或者一次性传输的数据块过大导致格式解析异常。小文件数据量小,缓冲区能正常处理,大文件就触发了这个问题。
方案一:调整ADBC批量写入参数
通过engine_kwargs给write_database传递batch_size参数,限制每次写入的数据行数,避免一次性传输过大的块:
def store_in_postgresql(df, table_name): password = 'anon' username = 'postgres' database = 'nyc_taxis' uri = f'postgresql://{username}:{password}@localhost:5432/{database}' common_sql_state = "SQLSTATE: 42P07" try: # 自定义批量写入大小,建议根据内存调整为5-20万行 df.write_database( table_name=table_name, connection=uri, engine='adbc', if_exists='replace', engine_kwargs={"batch_size": 100000} ) print('加载完成!') except Exception as e: err_str = str(e) if common_sql_state in err_str: df.write_database( table_name=table_name, connection=uri, engine='adbc', if_exists='append', engine_kwargs={"batch_size": 100000} ) print('加载完成!') else: print(e)
方案二:手动分片写入
如果调整参数无效,直接将DataFrame拆分成小分片,逐个写入数据库:
def store_in_postgresql(df, table_name, chunk_size=100000): password = 'anon' username = 'postgres' database = 'nyc_taxis' uri = f'postgresql://{username}:{password}@localhost:5432/{database}' common_sql_state = "SQLSTATE: 42P07" try: # 先尝试用第一个分片创建/替换表 first_chunk = df.slice(0, chunk_size) first_chunk.write_database( table_name=table_name, connection=uri, engine='adbc', if_exists='replace' ) # 循环写入剩余分片 for i in range(chunk_size, len(df), chunk_size): chunk = df.slice(i, chunk_size) chunk.write_database( table_name=table_name, connection=uri, engine='adbc', if_exists='append' ) print('加载完成!') except Exception as e: err_str = str(e) if common_sql_state in err_str: # 表已存在,直接追加所有分片 for i in range(0, len(df), chunk_size): chunk = df.slice(i, chunk_size) chunk.write_database( table_name=table_name, connection=uri, engine='adbc', if_exists='append' ) print('加载完成!') else: print(e)
调用时可根据实际情况调整chunk_size参数,平衡写入速度和内存占用。
额外优化建议
- 升级依赖到最新版本,旧版Polars或ADBC驱动可能存在大文件写入bug:
pip install --upgrade polars adbc-driver-postgresql - 检查PostgreSQL配置,适当调大
max_wal_size和work_mem参数,避免大文件写入时出现性能瓶颈或超时。
内容的提问来源于stack exchange,提问作者Zaid Allawanseh
相关产品推荐
相关产品推荐

