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

Python Polars处理大CSV写入PostgreSQL时内存泄漏崩溃求助

问题描述

我编写了一个Python Polars脚本,通过LazyFrame扫描CSV文件,选择指定列、过滤特定行后将结果上传至PostgreSQL数据库。该脚本在处理小型测试CSV时运行正常,但处理真实业务场景下的大CSV时,Python进程会持续占用内存直至崩溃。已在配备64GB内存的Windows设备和32GB内存的macOS设备上测试过,可排除硬件不足问题。尝试过设置low_memory=True、使用collect(engine='streaming')、sink_csv等方法,均无法解决内存占用过高的问题。

原代码如下:

import polars as pl
from sqlalchemy import create_engine

result = (
    pl.scan_csv("/path/to/file/201601.csv",
        separator="|",
        encoding='utf8-lossy',
        try_parse_dates=True
    )
    # 选择指定列
    .select(
        ["ProgramClassID", "NetworkAffiliationID", "LogServiceID", "LogEntryDate","StartTime", "EndTime","Duration",  "ProgramTitle", "Subtitle"]
    )
    # 过滤掉LogServiceID不等于3680的行
    .filter(  
        pl.col("LogServiceID") != 3680)  
    # 按日期排序
    .sort(  
        "LogEntryDate")  
    # 加载到内存
    .collect()  
)

# 数据库连接信息
uri='postgresql://postgres:pwd@0.0.0.0:5432/crtc_program_logs'

# 写入数据库
result.write_database(table_name="polars_test_02",
    connection=uri,
    if_table_exists="replace"
)

问题根源

核心问题出在排序操作和全量内存加载:

  1. LazyFrame的.sort()无法流式处理,必须将全量数据加载到内存才能完成排序,这是大文件内存溢出的主要原因。
  2. .collect()会把整个数据集加载到内存,后续write_database再全量写入,进一步加剧内存压力。

修复方案

1. 移除/转移排序操作

如果业务允许,直接去掉.sort();如果最终需要有序数据,建议在PostgreSQL中通过ORDER BY查询,或创建索引——数据库层面的排序比Python内存排序高效得多。

2. 流式写入数据库,避免全量加载

不要调用.collect(),直接在LazyFrame上调用.write_database(),Polars会自动分批次处理数据,无需将全量数据存入内存。

3. 优化CSV读取参数

  • 手动指定列的数据类型:通过dtypes参数给数值列指定合适的类型(如pl.Int32),减少Polars自动推断的内存开销。
  • 按需开启日期解析:如果try_parse_dates导致额外内存占用,可手动指定日期格式,或先按字符串读取再转换。

优化后的代码
import polars as pl

(
    pl.scan_csv(
        "/path/to/file/201601.csv",
        separator="|",
        encoding='utf8-lossy',
        try_parse_dates=True,
        # 手动指定数据类型,降低内存占用
        dtypes={
            "ProgramClassID": pl.Int32,
            "NetworkAffiliationID": pl.Int32,
            "LogServiceID": pl.Int32,
            "Duration": pl.Int32
        }
    )
    .select([
        "ProgramClassID", "NetworkAffiliationID", "LogServiceID", 
        "LogEntryDate","StartTime", "EndTime","Duration",  
        "ProgramTitle", "Subtitle"
    ])
    .filter(pl.col("LogServiceID") != 3680)
    # 直接流式写入数据库,无需加载全量数据
    .write_database(
        table_name="polars_test_02",
        connection='postgresql://postgres:pwd@0.0.0.0:5432/crtc_program_logs',
        if_table_exists="replace",
        # 控制批次大小,进一步降低内存峰值
        batch_size=100_000
    )
)

额外建议
  • 若必须在写入前排序:可按LogEntryDate范围拆分CSV,分别排序后写入数据库,最后在库内合并;或使用partition_by结合排序,但需严格控制单分区数据量。
  • 监控内存开销:通过pl.Config.set_tbl_rows(10)查看数据预览,或用profile()方法分析各步骤的内存占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.01 17:34:54