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

Polars合并多CSV时LazyFrame内存溢出,求高效解决方案

问题描述
  • 需要读取并合并4个CSV文件,文件列不完全一致但大部分重叠,缺列需填充NaN
  • 具体情况:
    1. 使用Polars将4个CSV读取到列表dfs中
    2. 每个CSV磁盘大小为2.8GB,df.estimated_size(unit='gb')返回约4.05GB
    3. 初始内存占用为31.2GB(来自htop)
    4. 尝试用pl.concat(dfs, parallel=False)合并
  • 结果:
    1. 使用scan_csv(延迟读取)时,128GB内存的VM崩溃
    2. 使用read_csv(急切读取)时,合并后内存从31.2GB升至50GB
  • 需求:读取并合并时避免内存重复占用,寻求Polars下内存高效的多CSV合并方法
现有代码
ND = 4
with open("./features_to_keep.pkl", "rb") as f:
    CTK = pickle.load(f)

CTK = set(CTK) #columns to keep. I will be reading in only a subset of these columns

#Schema 
with open("./obj_cols.pkl", "rb") as f:
    obj_cols = pickle.load(f)
    obj_cols = obj_cols.intersection(CTK)

with open("./bool_cols.pkl", "rb") as f:
    bool_cols = pickle.load(f)
    bool_cols = bool_cols.intersection(CTK)


def read_features(i):
    fp = f"features_{i}.csv"
    cols = pl.scan_csv(fp).columns
    cols = list(set(cols).intersection(CTK))
    dtypes = {i:(pl.datatypes.Categorical if i in obj_cols else pl.datatypes.Boolean if i in bool_cols else pl.datatypes.Float32) for i in cols}
    df = pl.read_csv(fp, columns = cols, dtypes = dtypes, infer_schema_length = 0)
    # df = pl.scan_csv(fp, dtypes = dtypes, infer_schema_length = 0).select(cols)
    return df


dfs = []
#I do not use stringcache context manager in case I had used scan_csv
with pl.StringCache():
    for i in range(ND):
        dfs.append(read_features(i))
X = pl.concat(dfs, parallel = False, how = 'diagonal')
内存高效解决方案

1. 正确使用延迟读取(LazyFrame)合并

原代码中scan_csv崩溃是因为调用.columns触发了全量文件扫描,且未在Lazy阶段完成合并。正确做法是直接构建LazyFrame列表,在Lazy层面完成合并后再collect(),Polars会自动优化执行计划,避免内存冗余:

ND = 4
with open("./features_to_keep.pkl", "rb") as f:
    CTK = pickle.load(f)

CTK = set(CTK)

with open("./obj_cols.pkl", "rb") as f:
    obj_cols = pickle.load(f)
    obj_cols = obj_cols.intersection(CTK)

with open("./bool_cols.pkl", "rb") as f:
    bool_cols = pickle.load(f)
    bool_cols = bool_cols.intersection(CTK)

def read_features_lazy(i):
    fp = f"features_{i}.csv"
    # 仅读取表头获取列名,避免全量扫描
    sample_df = pl.read_csv(fp, n_rows=0)
    cols = list(set(sample_df.columns).intersection(CTK))
    dtypes = {col: (pl.Categorical if col in obj_cols else pl.Boolean if col in bool_cols else pl.Float32) for col in cols}
    # 返回LazyFrame而非加载全量数据
    return pl.scan_csv(fp, columns=cols, dtypes=dtypes, infer_schema_length=0)

with pl.StringCache():
    lazy_dfs = [read_features_lazy(i) for i in range(ND)]
    # Lazy层面合并后再收集数据
    X = pl.concat(lazy_dfs, how='diagonal').collect()

2. 优化列读取逻辑

原代码用pl.scan_csv(fp).columns会触发全文件扫描,改用pl.read_csv(fp, n_rows=0)仅读取表头,速度更快且内存占用可以忽略。

3. 分块写入磁盘(极端内存压力下可选)

如果Lazy合并仍有内存压力,可以将每个文件的内容分块写入Parquet文件,避免全量数据驻留内存:

with pl.StringCache():
    # 初始化第一个文件的Parquet写入
    first_lazy = read_features_lazy(0)
    first_lazy.sink_parquet("merged_data.parquet")
    # 后续文件追加写入
    for i in range(1, ND):
        current_lazy = read_features_lazy(i)
        current_lazy.sink_parquet("merged_data.parquet", append=True)
    # 最后读取合并后的文件
    X = pl.read_parquet("merged_data.parquet")

4. 确认数据类型优化

确保Float32、Boolean、Categorical类型正确应用:

  • Boolean类型每列仅占1bit/行,比默认整数类型节省大量内存
  • Categorical对重复字符串列的压缩效果显著
  • Float32比Float64内存占用减半,多数场景下精度足够

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 21:42:40