从Pickle文件创建大型Polars DataFrame时如何限制内存占用(上限15GB)
限制Pickle合并脚本内存至15GB以内的实现方法
核心优化方向
原脚本通过Pandas中转读取Pickle再逐步合并,容易造成内存累积。以下是几种可落地的优化方案:
1. 跳过Pandas中转,直接用Polars读取Pickle
Polars原生支持读取Pickle文件,无需通过Pandas转换,能减少临时内存占用。
修改后的代码:
import polars as pl import glob pickle_files = glob.glob("/home/x/pickles/*.pkl.gz") df_polars = pl.DataFrame() for file in pickle_files: # 直接用Polars读取,避免Pandas中转的额外内存开销 df_temp = pl.read_pickle(file) df_polars = df_polars.vstack(df_temp) # 手动释放临时变量内存 del df_temp print(df_polars)
2. 分批合并+内存释放
如果单文件内存占用较高,可计算单次可加载的文件数量,分批次合并后写入Parquet临时文件(Parquet比Pickle更省空间),最后再合并所有临时文件,全程控制内存不超15GB。
示例代码:
import polars as pl import glob import gc import os pickle_files = glob.glob("/home/x/pickles/*.pkl.gz") # 根据单文件内存调整批次大小,比如单文件1GB就设为15 batch_size = 10 temp_files = [] # 分批次处理文件 for i in range(0, len(pickle_files), batch_size): batch_files = pickle_files[i:i+batch_size] batch_df = pl.concat([pl.read_pickle(f) for f in batch_files]) temp_path = f"/tmp/temp_batch_{i}.parquet" batch_df.write_parquet(temp_path) temp_files.append(temp_path) # 释放当前批次内存 del batch_df gc.collect() # 合并所有临时文件 final_df = pl.concat([pl.read_parquet(f) for f in temp_files]) print(final_df) # 清理临时文件(可选) for f in temp_files: os.remove(f)
3. 使用Lazy模式延迟加载(适合数据分析场景)
如果最终目的是做数据分析而非一次性持有全量数据,用Polars的LazyFrame可以完全避免加载全量数据到内存,按需执行计算。
示例代码:
import polars as pl import glob pickle_files = glob.glob("/home/x/pickles/*.pkl.gz") # 创建LazyFrame集合,仅记录读取逻辑不加载数据 lazy_dfs = [pl.scan_pickle(file) for file in pickle_files] # 合并为一个LazyFrame final_lazy_df = pl.concat(lazy_dfs) # 后续操作直接基于LazyFrame执行,比如过滤、聚合,仅在collect()时加载必要数据 result = final_lazy_df.filter(pl.col("some_column") > 100)\ .group_by("category")\ .agg(pl.col("value").sum())\ .collect() print(result)
4. 实时监控内存并动态调整
加入内存监控,实时判断合并后内存是否超阈值,超限时写入临时文件清空当前数据,避免内存溢出。
示例代码:
import polars as pl import glob import psutil import os def get_current_memory_usage(): process = psutil.Process() # 返回当前进程内存占用(GB) return process.memory_info().rss / (1024**3) pickle_files = glob.glob("/home/x/pickles/*.pkl.gz") df_polars = pl.DataFrame() max_memory = 15 # 内存上限(GB) temp_overflow_path = "/tmp/temp_overflow.parquet" for file in pickle_files: df_temp = pl.read_pickle(file) # 预估合并后的内存占用 estimated_new_size = (df_polars.estimated_size() + df_temp.estimated_size()) / (1024**3) if estimated_new_size > max_memory: # 超过阈值则写入临时文件,清空当前DataFrame df_polars.write_parquet(temp_overflow_path) df_polars = pl.DataFrame() df_polars = df_polars.vstack(df_temp) del df_temp # 合并最后剩余数据与临时文件 if os.path.exists(temp_overflow_path): temp_df = pl.read_parquet(temp_overflow_path) df_polars = pl.concat([temp_df, df_polars]) os.remove(temp_overflow_path) print(df_polars)
内容的提问来源于stack exchange,提问作者PaulS
相关产品推荐
相关产品推荐

