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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 10:01:34