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

如何利用40核120GB内存高效聚合600GB CSV数据?

大规模CSV聚合的内存处理方案(基于40核CPU+120GB RAM)

核心结论

直接用Pandas全量加载处理完全不可行;Dask是最适配你场景的解决方案;Modin可行但有明显局限性;「读取全部CSV到数据框再聚合」的思路完全不可行。


一、为什么Pandas直接处理行不通

Pandas是单进程模型,要求数据全量加载到内存。你的600GB原始数据远超过120GB RAM,即使CSV有压缩,解压后体积只会更大,直接读取必然触发内存溢出。另外,8000万个小文件的IO开销会让加载过程慢到无法接受。


二、Dask:最优选择

Dask专为大规模数据的并行分块处理设计,完美适配你的硬件资源:

  • 解决小文件IO瓶颈:8000万个小文件会带来极高的元数据读取延迟,建议先批量合并成GB级大文件:
    import dask.bag as db
    
    # 读取所有小文件
    file_bag = db.read_text("your_data_dir/*.csv")
    # 按10GB为单位合并成大文件
    file_bag.map_partitions(lambda chunk: '\n'.join(chunk)).to_textfiles("merged_dir/part-*.csv", blocksize="10GB")
    
  • 内存可控的并行读取:用Dask DataFrame读取合并后的文件,通过blocksize控制单个分区大小(比如设为2GB),确保每个分区能被内存容纳,同时利用40核并行处理多个分区:
    import dask.dataframe as dd
    
    df = dd.read_csv("merged_dir/*.csv", blocksize="2GB", usecols=["需要的字段1", "需要的字段2"])
    
  • 兼容Pandas的聚合操作:Dask DataFrame支持和Pandas几乎一致的聚合API,自动拆分任务到多核并行执行,最后合并结果:
    # 示例:按network_id分组统计流量总和
    agg_result = df.groupby("network_id").agg({"traffic": "sum"}).compute()
    # 保存聚合结果
    agg_result.to_csv("aggregated_result.csv", index=False)
    

三、Modin:可行但有局限

Modin是Pandas的并行封装,能利用多核,但并不适配你的场景:

  • 对8000万个小文件的IO调度效率远低于Dask,元数据加载会非常缓慢;
  • 当数据总量远超内存时,它依赖底层Ray/Dask引擎做内存溢出处理,但性能和灵活性不如原生Dask;
  • 仅在已经合并小文件、且聚合逻辑简单的前提下可以尝试,整体优先级低于Dask。

四、"全量读取到数据框再聚合"不可行的原因

  1. 内存不足:600GB数据远超过120GB RAM,全量加载必然触发OOM;
  2. IO开销过大:8000万个小文件的读取延迟会让加载过程耗时极长,根本无法完成全量加载。

额外优化建议

  • 字段过滤:读取时用usecols指定仅加载聚合所需的字段,减少内存占用;
  • 类型优化:提前指定列的数据类型(比如把字符串列设为category,数值列用int32/float32替代默认的int64/float64),进一步压缩内存开销;
  • 分区数匹配:设置npartitions等于或略大于CPU核心数(比如40或80),充分利用多核并行能力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 15:00:25