使用Dask合并超5GB大CSV文件导出时内存不足的技术咨询
我之前处理大文件合并时也碰到过一模一样的情况——8GB内存面对几个5GB+的CSV,合并后导出分分钟爆内存。咱们从几个关键点入手优化,应该就能搞定:
1. 先搞清楚内存爆掉的原因
Dask虽然是懒加载,但导出时如果合并后的分区太大,单个分区加载到内存时就会超出你的8GB限制;另外你把所有列的dtype设为object,这会比用具体类型(比如int、float、string)占用多得多的内存,也是隐形的内存杀手。
2. 针对性优化方案
优化一:指定具体数据类型,砍掉不必要的内存占用
别再全用object了!先大概看下你的CSV里各列的类型,然后给read_csv传一个dtype字典,比如数值列用int64/float64,字符串列用string,日期列用datetime64[ns]。这一步能大幅降低内存占用:
dtype_dict = { "A": "int64", # 假设A是ID列,整数类型 "B": "float64", # 数值列 "C": "string", # 字符串列 # 剩下的列根据实际情况补充 } data_1 = dd.read_csv(file_loc_1, dtype=dtype_dict, encoding='cp1252') data_2 = dd.read_csv(file_loc_2, dtype=dtype_dict, encoding='cp1252')
优化二:合并后重新分区,让每个分区“轻量”起来
合并后的DataFrame分区可能很大,导出时单个分区加载到内存就会触发内存不足。用repartition把分区拆小,比如每个分区控制在100MB左右(可以根据你的内存调整,比如50MB):
final_1 = dd.merge(data_1, data_2, left_on="A", right_on="A", how="left") # 按大小分区,每个分区约100MB final_1 = final_1.repartition(partition_size='100MB')
优化三:导出时避开“一次性加载全量数据”的坑
如果直接导出单个大文件,Dask可能会尝试把所有分区合并后再写入,这很容易爆内存。推荐两种方式:
方式A:先导出多个小文件,再用系统命令合并
这种方式最稳妥,系统级别的合并不需要加载整个文件到内存:
# 导出为多个临时文件,命名格式是temp_output_00.csv、temp_output_01.csv... final_1.to_csv('temp_output_*.csv', encoding='cp1252', index=False)
然后根据你的操作系统用命令合并:
- Windows:打开命令提示符,执行
copy temp_output_*.csv final_merged.csv - Linux/macOS:执行
cat temp_output_*.csv > final_merged.csv
方式B:强制导出单个文件(需配合小分区)
如果一定要用Dask直接导出单个文件,确保分区已经足够小,然后用single_file=True参数:
final_1.to_csv('final_merged.csv', single_file=True, encoding='cp1252', index=False)
Dask会逐个处理分区并追加到同一个文件,不会一次性加载所有数据。
优化四:限制Dask的内存使用
给Dask设置内存上限,避免它占满你的8GB内存导致系统崩溃:
from dask.distributed import Client # 给Dask分配6GB内存,留2GB给系统其他进程 client = Client(memory_limit='6GB')
3. 最终修改后的完整代码示例
import dask.dataframe as dd from dask.distributed import Client # 初始化Client,限制内存使用 client = Client(memory_limit='6GB') # 定义具体的数据类型,减少内存占用 dtype_dict = { "A": "int64", "B": "float64", "C": "string", # 根据你的实际列补充 } file_loc_1 = r"..." file_loc_2 = r"..." data_1 = dd.read_csv(file_loc_1, dtype=dtype_dict, encoding='cp1252') data_2 = dd.read_csv(file_loc_2, dtype=dtype_dict, encoding='cp1252') # 执行左连接 final_1 = dd.merge(data_1, data_2, left_on="A", right_on="A", how="left") # 重新分区,减小单个分区大小 final_1 = final_1.repartition(partition_size='100MB') # 导出为多个临时文件 final_1.to_csv('temp_output_*.csv', encoding='cp1252', index=False) # 后续手动用系统命令合并临时文件即可
如果还是有问题,可以试试把分区大小再调小(比如50MB),或者先删除不需要的列,进一步减少数据量。
内容的提问来源于stack exchange,提问作者ivyyi

