如何利用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。
四、"全量读取到数据框再聚合"不可行的原因
- 内存不足:600GB数据远超过120GB RAM,全量加载必然触发OOM;
- IO开销过大:8000万个小文件的读取延迟会让加载过程耗时极长,根本无法完成全量加载。
额外优化建议
- 字段过滤:读取时用
usecols指定仅加载聚合所需的字段,减少内存占用; - 类型优化:提前指定列的数据类型(比如把字符串列设为
category,数值列用int32/float32替代默认的int64/float64),进一步压缩内存开销; - 分区数匹配:设置
npartitions等于或略大于CPU核心数(比如40或80),充分利用多核并行能力。
内容的提问来源于stack exchange,提问作者Ranger
相关产品推荐
相关产品推荐

