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

使用Polars合并数千个CSV/Feather文件时内核崩溃的解决方法

解决Polars合并数千个CSV/Feather文件时内核崩溃的问题

问题根源

你当前的代码存在两个核心问题:

  1. 内存累积:collected_temp_all_spatial_join_data是驻留内存的DataFrame,每次拼接都会持续占用更多内存,数千个文件处理后会直接耗尽内存。
  2. 查询计划膨胀:反复拼接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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 22:42:54