使用Polars合并数千个CSV/Feather文件时内核崩溃的解决方法
解决Polars合并数千个CSV/Feather文件时内核崩溃的问题
问题根源
你当前的代码存在两个核心问题:
- 内存累积:
collected_temp_all_spatial_join_data是驻留内存的DataFrame,每次拼接都会持续占用更多内存,数千个文件处理后会直接耗尽内存。 - 查询计划膨胀:反复拼接LazyFrame会让Polars维护的查询计划变得异常庞大,最终超出内核的内存承载能力,导致崩溃。
优化方案
方案1:利用Polars懒加载特性直接批量处理(推荐)
Polars支持直接批量扫描所有文件并生成统一的LazyFrame,写入Parquet时会自动并行处理、分块读写,内存占用可控,完全不需要手动分批。
import os import glob import polars as pl # 替换为你的实际路径和参数 EXPORT_TABLE_FOLDER = "你的目标文件夹路径" location_name = "目标位置名称" table_format = "csv" # 可选值:"csv" / "feather" parquet_name = f"simplified_all_spatial_join_data_{location_name}_p0.parquet" parquet_path = os.path.join(EXPORT_TABLE_FOLDER, parquet_name) if not os.path.exists(parquet_path): # 获取所有目标文件路径 all_tables = glob.glob(os.path.join(EXPORT_TABLE_FOLDER, f"*.{table_format}")) # 根据文件格式选择扫描函数 scan_func = pl.scan_csv if table_format == "csv" else pl.scan_ipc # 批量扫描所有文件生成统一LazyFrame combined_lf = pl.concat([scan_func(path, infer_schema_length=0) for path in all_tables], how="vertical") # 直接写入Parquet,Polars自动处理内存和并行 combined_lf.write_parquet(parquet_path) else: print("WARNING: 文件已存在,跳过处理流程")
方案2:分批写入临时文件再合并(极端内存受限场景)
如果你的内存不足以支撑一次性扫描所有文件的查询计划,可以分批处理并写入临时Parquet,最后合并所有临时文件:
import os import glob import polars as pl from tqdm import tqdm # 替换为你的实际路径和参数 EXPORT_TABLE_FOLDER = "你的目标文件夹路径" location_name = "目标位置名称" table_format = "csv" # 可选值:"csv" / "feather" batch_size = 50 # 每批处理的文件数量,可根据内存调整 parquet_name = f"simplified_all_spatial_join_data_{location_name}_p0.parquet" parquet_path = os.path.join(EXPORT_TABLE_FOLDER, parquet_name) temp_folder = os.path.join(EXPORT_TABLE_FOLDER, "temp_parquet_batches") os.makedirs(temp_folder, exist_ok=True) if not os.path.exists(parquet_path): all_tables = glob.glob(os.path.join(EXPORT_TABLE_FOLDER, f"*.{table_format}")) scan_func = pl.scan_csv if table_format == "csv" else pl.scan_ipc # 分批处理并写入临时文件 for batch_idx in tqdm(range(0, len(all_tables), batch_size)): batch_files = all_tables[batch_idx:batch_idx+batch_size] batch_lf = pl.concat([scan_func(path, infer_schema_length=0) for path in batch_files], how="vertical") temp_parquet_path = os.path.join(temp_folder, f"batch_{batch_idx}.parquet") batch_lf.write_parquet(temp_parquet_path) # 合并所有临时Parquet文件 temp_parquets = glob.glob(os.path.join(temp_folder, "*.parquet")) final_lf = pl.concat([pl.scan_parquet(path) for path in temp_parquets], how="vertical") final_lf.write_parquet(parquet_path) # 清理临时文件(可选) for temp_file in temp_parquets: os.remove(temp_file) os.rmdir(temp_folder) else: print("WARNING: 文件已存在,跳过处理流程")
原代码问题总结
- 手动在内存中累积DataFrame会导致内存持续增长,最终耗尽资源。
- 反复拼接LazyFrame会让查询计划无限膨胀,超出内核的内存管理能力。
内容的提问来源于stack exchange,提问作者Amri Rasyidi
相关产品推荐
相关产品推荐

