导出千万级Dask DataFrame至TXT失败:报错修复与性能优化咨询
问题修复与代码优化方案
一、pynvml错误的核心原因及修复
原代码的逐行compute()导出方式会导致内存暴增,触发系统资源耗尽,进而引发pynvml(NVIDIA设备管理库)的异常。修复该错误的核心是改变导出逻辑,避免逐行加载计算,同时优化内存占用。若确实涉及GPU资源冲突,可强制Dask使用CPU计算:
import os os.environ["CUDA_VISIBLE_DEVICES"] = "" # 禁用GPU,强制CPU计算
二、代码优化全方案
1. 读取阶段优化:减少中间内存开销
直接在dd.read_csv中指定所需列和数据类型,避免后续drop和astype操作,同时批量读取所有CSV文件,利用Dask并行能力:
import glob import re import dask.dataframe as dd import os csv_files = glob.glob("xxx_*.csv") used_cols = ["word", "word_freq", "doc_freq", "advis_word_freq", "advis_doc_freq", "story_word_freq", "story_doc_freq", "multi_word_freq", "multi_doc_freq", "other_word_freq", "other_doc_freq"] # 提前定义数据类型,减少内存占用 dtype_spec = {col: "int16" for col in used_cols if col != "word"} pattern = re.compile('[一-鿿]+_[一-鿿]+_[一-鿿]+') # 批量读取所有CSV,直接过滤列和指定类型 df = dd.read_csv( csv_files, encoding="UTF-8", usecols=used_cols, dtype=dtype_spec ) # 过滤符合规则的word行 df = df[df['word'].str.contains(pattern)]
2. 聚合与分区优化:平衡计算效率与内存
分组聚合后,合理设置分区数(建议每个分区控制在100MB左右,1000万行数据设置20-50个分区即可,过多分区会增加调度开销):
result = df.groupby("word").sum().reset_index() # 调整分区数,根据自身内存情况灵活修改 result = result.repartition(npartitions=30) print(f"result npartition = {result.npartitions}")
3. 导出TXT优化:高效批量写入
使用map_partitions并行处理每个分区,以追加模式写入文件,彻底避免逐行计算的低效与内存问题:
def write_partition(partition, output_path): # 将分区每行转为逗号分隔的字符串 lines = partition.apply(lambda row: ", ".join(map(str, row)), axis=1) with open(output_path, "a", encoding="UTF-8") as f: f.write("\n".join(lines) + "\n") return None # 先清空目标文件 with open("aaa.txt", "w", encoding="UTF-8"): pass # 并行写入所有分区 result.map_partitions(write_partition, "aaa.txt").compute() print("Done!")
替代导出方案:先存临时CSV再合并
如果担心map_partitions的并发写入冲突,可先导出为多个临时CSV,再合并为单个TXT:
# 导出为临时CSV(无表头、无索引) result.to_csv("temp_export_*.csv", encoding="UTF-8", index=False, header=False) # 合并所有临时文件到aaa.txt with open("aaa.txt", "wb") as out_file: for temp_file in glob.glob("temp_export_*.csv"): with open(temp_file, "rb") as in_file: out_file.write(in_file.read()) os.remove(temp_file) # 删除临时文件
三、关键优化点总结
- 避免手动循环读取CSV,利用Dask批量读取的并行能力提升效率。
- 读取阶段直接指定列和数据类型,减少内存占用与后续转换开销。
- 绝对禁止逐行
compute(),改用分区级批量处理,大幅降低内存压力。 - 合理设置分区数,平衡计算并行度与调度开销。
内容的提问来源于stack exchange,提问作者Kevin.M.t
相关产品推荐
相关产品推荐

