如何在Python中处理2000万行CSV并删除client_id列重复项
解决大型CSV去重(删除client_id重复行)的问题
你的Dask代码里的关键问题
- API误用+拼写错误:
drop_duplicates里的Inplace=True完全错误——首先参数名应为小写的inplace,其次Dask DataFrame不支持inplace操作(它是不可变结构),必须用返回的新对象重新赋值。 - 全量读取检测编码:
chardet.detect(f.read())会把整个大文件塞进内存,对超大CSV来说直接导致内存溢出或巨慢,只需读取前几KB就能完成编码检测。 - 直接覆盖原文件:
to_csv('bigData.csv')会在读取原文件的同时写入,大概率导致文件损坏或读取异常,应该先写临时文件再替换。 - 重复计算浪费资源:两次调用
compute()会触发两次全量数据扫描,耗时直接翻倍。 - 未启用进度条:导入了
ProgressBar但没激活,看不到执行进度,误以为程序卡住。
修正后的Dask代码
import dask.dataframe as dd import chardet from dask.diagnostics import ProgressBar import shutil import os # 仅读取前10KB检测编码,避免加载全量数据 with open('bigData.csv', 'rb') as f: result = chardet.detect(f.read(10240)) # 10KB足够完成编码检测 # 激活进度条,实时查看执行状态 ProgressBar().register() # 读取CSV,指定分隔符与编码 df = dd.read_csv('bigData.csv', encoding=result['encoding'], sep=';') # 一次compute获取原始行数和去重后行数,避免重复扫描 total_rows, deduped_rows = dd.compute( df.shape[0], df.drop_duplicates(subset=['client_id'], keep=False).shape[0] ) # 执行去重并写入临时单文件 deduped_df = df.drop_duplicates(subset=['client_id'], keep=False) deduped_df.to_csv('bigData_temp.csv', sep=';', index=False, single_file=True) # 替换原文件(若需保留原文件,可跳过此步直接使用temp文件) shutil.move('bigData_temp.csv', 'bigData.csv') # 输出结果 total_duplicates = total_rows - deduped_rows print(f'删除了 {total_duplicates} 条重复行。')
额外优化方案
如果文件大到Dask处理仍缓慢,试试命令行工具awk——速度更快且内存占用极低:
# 第一步:统计所有client_id的出现次数,记录仅出现一次的ID awk -F ';' '{count[$1]++} END {for (id in count) if (count[id]==1) print id}' bigData.csv > keep_ids.txt # 第二步:保留表头和符合条件的行 awk -F ';' 'NR==1 || $1 in keep' bigData.csv keep_ids.txt > bigData_deduped.csv
- 调整Dask分区:在
read_csv中添加blocksize='64MB'(根据你的内存情况调整),让分区大小匹配硬件能力,提升处理效率。 - 彻底放弃Pandas:Pandas会把全量数据加载到内存,内存不足时必然卡顿崩溃,Dask是处理超大文件的正确选择,但要避免上述低级错误。
内容的提问来源于stack exchange,提问作者David Kopl
相关产品推荐
相关产品推荐

