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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 10:12:05