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

使用Dask处理文本后调用to_csv()遇内存问题求解决方案

Dask处理文本数据后保存CSV的内存问题解决方案

问题背景

清洗大文本数据集后,调用Dask的to_csv()方法时遭遇内存溢出问题。已尝试调整worker数量、blocksize、内存限制,改用apply()、persist()及Pandas分块写入等方法,均未解决问题。

实用解决技巧

  • 精准控制分区写入节奏
    避免Dask一次性处理所有分区写入,手动逐个执行分区任务,降低瞬时内存占用:
    # 先生成写入任务,不立即执行
    write_tasks = df.to_csv('cleaned_lemma_*.csv', index=False, encoding='utf-8', single_file=False, compute=False)
    # 逐个执行分区写入
    for task in write_tasks:
        task.compute()
    
  • 缩小分区尺寸
    进一步降低blocksize(比如设为10e6即10MB),或手动重新拆分现有分区,确保单个分区内存占用不超过worker内存的1/3:
    # 将现有分区拆分为原数量的2倍,减小单分区大小
    df = df.repartition(npartitions=df.npartitions * 2)
    df.to_csv('cleaned_lemma_*.csv', index=False, encoding='utf-8', single_file=False)
    
  • 彻底清理缓存
    在写入前主动清理Dask的计算缓存与Python垃圾,释放闲置内存:
    # 取消Dask中未完成的关联任务缓存
    client.cancel(df)
    # 强制触发Python垃圾回收
    gc.collect()
    # 执行写入
    df.to_csv('cleaned_lemma_*.csv', index=False, encoding='utf-8', single_file=False)
    
  • 排查清洗函数内存泄漏
    检查wrapper_func是否存在内存泄漏(比如创建全局大对象、未释放临时变量),可在函数中添加内存监控:
    import psutil
    def wrapper_with_mem_check(text):
        cleaned_text = wrapper_func(text)
        # 打印当前分区处理后的内存占用
        print(f"Partition memory usage: {psutil.Process().memory_info().rss / 1024**2:.2f} MB")
        return cleaned_text
    df['cleaned_title'] = df['title'].map_partitions(lambda p: p.apply(wrapper_with_mem_check), meta=('title', 'object'))
    

原始代码

import time
import gc
from dask.distributed import Client
import dask.dataframe as dd

if __name__ == '__main__':
    # Initialize the Dask client
    client = Client(n_workers=3, threads_per_worker=1, memory_limit='1.5GB')
    print('Dask client created')

    PATH = "C:\\Users\\el ruchenzo\\jobsproject\\jobsproject\\lt_data.csv"
    reqd = ['description', 'title', 'code']
    blocksize = 25e6  # REDUCED FROM 100 GB TO 25 GB

    # Load the CSV with Dask
    df = dd.read_csv(PATH,
                     usecols=reqd,
                     blocksize=blocksize,
                     dtype={'Code': 'float'},
                     engine='python',
                     encoding='utf-8',
                     on_bad_lines='skip')

    # Apply the cleaning function to the 'title' column in the DataFrame
    start_time = time.time()
    df['cleaned_title'] = df['title'].map_partitions(lambda partition: partition.apply(wrapper_func), meta=('title', 'object'))
    gc.collect()
    end_time = time.time()
    print(f"1. Processing time: {end_time - start_time:.2f} seconds")

    # Apply the cleaning function to the 'description' column in the DataFrame
    start_time = time.time()
    df['cleaned_description'] = df['description'].map_partitions(lambda partition: partition.apply(wrapper_func), meta=('description', 'object'))
    gc.collect()
    end_time = time.time()
    print(f"2. Processing time: {end_time - start_time:.2f} seconds")

    df.to_csv('cleaned_lemma_*.csv', index=False, encoding='utf-8', single_file = False)
    print('Saved to csv successfully')

    print('Work ended successfully')

内容的提问来源于stack exchange,提问作者Rūta Kondrot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 09:04:55