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

如何清除polars.read_csv()读取后占用的RSS内存?

问题描述

一次性读取大型CSV文件负载过高,因此尝试分批处理:将文件物理拆分为10份。理论上无论读取单份还是循环读取10次,内存占用都不应持续增长,但实际却不断上升。
直接创建DataFrame对象而非使用read_csv时,内存不会增长,因此推测read_csv存在内存泄漏。使用gc.collect()或del关键字均无法解决该问题,寻求更好的处理方案。

使用版本

polars==1.11.0

测试代码

import polars as pl

import multiprocessing as mp
import psutil
import os
import gc
import glob
import sys


def read_one():
    mypid = os.getpid()
    proc = psutil.Process(mypid)

    columns = ['a', 'b', 'c']

    filepath = 'test1.csv'

    current_rss = proc.memory_info().rss

    df = pl.read_csv(filepath, has_header=False, new_columns=columns, schema_overrides={c: pl.String for c in columns})

    current_rss2 = proc.memory_info().rss

    print('rss size:', current_rss2-current_rss)


def read_all():
    mypid = os.getpid()
    proc = psutil.Process(mypid)

    columns = ['a', 'b', 'c']

    filepath = 'test*.csv'

    current_rss = proc.memory_info().rss

    df = pl.read_csv(filepath, has_header=False, new_columns=columns, schema_overrides={c: pl.String for c in columns})

    current_rss2 = proc.memory_info().rss

    print('rss size:', current_rss2-current_rss)

def read_loop():
    mypid = os.getpid()
    proc = psutil.Process(mypid)

    columns = ['a', 'b', 'c']

    filelist = glob.glob('test*.csv')

    current_rss = proc.memory_info().rss

    for filepath in filelist:
        df = pl.read_csv(filepath, has_header=False, new_columns=columns,
                         schema_overrides={c: pl.String for c in columns})

        del df
        gc.collect()

    current_rss2 = proc.memory_info().rss

    print('rss size:', current_rss2 - current_rss)


if __name__ == '__main__':
    mp.set_start_method('spawn')

    proc = mp.Process(target=read_one, daemon=True)
    proc.start()
    proc.join()

    proc = mp.Process(target=read_all, daemon=True)
    proc.start()
    proc.join()

    proc = mp.Process(target=read_loop, daemon=True)
    proc.start()
    proc.join()

输出结果

rss size: 213282816
rss size: 2039447552
rss size: 423301120
解决方案
  • 升级Polars版本:Polars 1.11.0属于旧版本,后续稳定版(如1.14+)已修复多个内存泄漏相关Bug,执行以下命令升级:

    pip install --upgrade polars
    
  • 使用流式读取替代循环单文件读取:通过scan_csv实现延迟加载,分批处理数据,避免一次性加载全部内容到内存:

    def read_stream():
        mypid = os.getpid()
        proc = psutil.Process(mypid)
        columns = ['a', 'b', 'c']
        current_rss = proc.memory_info().rss
        
        # 流式扫描所有CSV文件
        lf = pl.scan_csv('test*.csv', has_header=False, new_columns=columns, schema_overrides={c: pl.String for c in columns})
        # 按批次处理,示例为每次处理10000行
        batch_size = 10000
        for batch in lf.iter_batches(batch_size=batch_size):
            # 在此处添加你的数据处理逻辑
            pass
        
        current_rss2 = proc.memory_info().rss
        print('rss size:', current_rss2 - current_rss)
    
  • 手动清除Polars内部缓存:Polars的C++内核可能存在未及时释放的内部缓存,每次读取后调用pl.clear_cache()辅助释放:

    def read_loop_improved():
        mypid = os.getpid()
        proc = psutil.Process(mypid)
        columns = ['a', 'b', 'c']
        filelist = glob.glob('test*.csv')
        current_rss = proc.memory_info().rss
    
        for filepath in filelist:
            df = pl.read_csv(filepath, has_header=False, new_columns=columns, schema_overrides={c: pl.String for c in columns})
            del df
            gc.collect()
            pl.clear_cache()  # 清除Polars内部缓存
    
        current_rss2 = proc.memory_info().rss
        print('rss size:', current_rss2 - current_rss)
    
  • 多进程隔离单次读取:将每个文件的读取和处理放在独立子进程中,子进程结束后会自动释放所有内存,避免主进程内存累积:

    def process_file(filepath, columns):
        df = pl.read_csv(filepath, has_header=False, new_columns=columns, schema_overrides={c: pl.String for c in columns})
        # 在此处添加你的数据处理逻辑
        return
    
    def read_multi_process():
        mypid = os.getpid()
        proc = psutil.Process(mypid)
        columns = ['a', 'b', 'c']
        filelist = glob.glob('test*.csv')
        current_rss = proc.memory_info().rss
    
        for filepath in filelist:
            proc = mp.Process(target=process_file, args=(filepath, columns), daemon=True)
            proc.start()
            proc.join()
    
        current_rss2 = proc.memory_info().rss
        print('rss size:', current_rss2 - current_rss)
    

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 02:47:26