使用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
相关产品推荐
相关产品推荐

