如何清除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
相关产品推荐
相关产品推荐

