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

导出千万级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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 17:15:35