Polars合并多CSV时LazyFrame内存溢出,求高效解决方案
问题描述
- 需要读取并合并4个CSV文件,文件列不完全一致但大部分重叠,缺列需填充NaN
- 具体情况:
- 使用Polars将4个CSV读取到列表
dfs中 - 每个CSV磁盘大小为2.8GB,
df.estimated_size(unit='gb')返回约4.05GB - 初始内存占用为31.2GB(来自htop)
- 尝试用
pl.concat(dfs, parallel=False)合并
- 使用Polars将4个CSV读取到列表
- 结果:
- 使用
scan_csv(延迟读取)时,128GB内存的VM崩溃 - 使用
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
相关产品推荐
相关产品推荐

