pandas中使用Levenshtein做字符串比对的Python代码性能优化
问题根因
代码运行慢、内存溢出和用pandas还是modin没有关系,核心是两个硬伤:
- 算法复杂度达到了O(n²):n行文本就要做n*n次两两比对,1万行数据就要生成1亿个相似度值,光存储这个相似度矩阵就要占用数GB内存,计算量随数据量平方级上涨,数据量翻10倍计算量涨100倍,跑几天、爆内存是必然结果
- 冗余操作太多:先存全量相似度矩阵再逐行逐列循环把上三角置0,纯Python逐元素遍历的开销极高;modin对这种嵌套全量遍历的自定义逻辑没有优化效果,反而会因为分布式数据拷贝增加额外开销。
可落地优化方案
1. 前置剪枝,砍掉99%无意义的比对
fuzz.partial_ratio要两个文本存在足够长度的公共子串才会拿到≥75的分数,根本不需要两两全量计算:
- 对所有text做3-gram切分,建立倒排索引,记录每个gram出现在哪些文本里
- 实际计算相似度前,只和共享至少2个公共gram的文本做fuzz计算,其余文本直接判定为相似度低于阈值,跳过计算
- 这一步能把需要实际计算fuzz的配对数从O(n²)降到接近O(n)级别,性能提升可达数百到数千倍
2. 替换计算库,砍掉冗余内存占用
- 把
thefuzz换成rapidfuzz,两者接口完全兼容,rapidfuzz是C++实现,相同逻辑计算速度是thefuzz的5-10倍,内存占用更低 - 完全放弃存储n*n全量相似度矩阵的逻辑:逐行判定时,只和索引大于当前行的候选文本比对,只要碰到一个相似度≥75的配对,直接标记当前行为需要过滤,立刻终止当前行的后续计算,不需要算完所有配对。举个例子,某行文本和第3个候选的相似度就达到80分,后面剩下的几千、几万次比对全部可以跳过,能省掉巨量无效计算。
3. 用原生多进程替代modin,充分利用多核性能
modin对你这种自定义遍历逻辑的优化效果极差,直接用Python原生multiprocessing实现并行:
- 把待判定的文本按行号切分成和CPU核心数匹配的块,每个进程负责一个块的相似度判定
- 进程间只传递必要的文本列表和索引,避免大DataFrame的跨进程拷贝开销,32核机器开28-30个进程,计算速度可以线性提升近30倍
4. 大文件按需读入,降低内存峰值
- 读CSV时先用
usecols参数只加载需要用到的列,不要一开始就把30个字段全读进内存 - 单文件过大时用
chunksize参数按块读入,块内先做初步过滤,再和已判定为保留的文本做比对,避免一次性加载全表占满内存。
优化后参考代码
import glob import pandas as pd from rapidfuzz import fuzz from multiprocessing import Pool, cpu_count from collections import defaultdict # 构建3-gram倒排索引 def build_inverted_index(texts, ngram_size=3): index = defaultdict(set) for idx, text in enumerate(texts): if not isinstance(text, str): text = "" if len(text) < ngram_size: index[text].add(idx) continue for i in range(len(text) - ngram_size + 1): gram = text[i:i+ngram_size] index[gram].add(idx) return index # 单块文本判定:存在任意其他文本相似度≥75就标记为过滤 def check_chunk(args): chunk_rows, chunk_start_idx, all_texts, inv_index, threshold, min_common_gram = args keep_flag = [] for row_offset, text in enumerate(chunk_rows): current_idx = chunk_start_idx + row_offset if not isinstance(text, str): keep_flag.append(True) continue # 召回共享gram的候选 candidates = set() text_len = len(text) if text_len < 3: grams = [text] else: grams = [text[i:i+3] for i in range(text_len-2)] for g in grams: candidates.update(inv_index.get(g, set())) # 只和索引更大的候选比对,避免重复计算 is_duplicate = False for cand_idx in candidates: if cand_idx <= current_idx: continue # 长度差过大直接跳过 if len(all_texts[cand_idx]) < text_len * 0.6: continue score = fuzz.partial_ratio(text, all_texts[cand_idx], score_cutoff=threshold) if score >= threshold: is_duplicate = True break keep_flag.append(not is_duplicate) return keep_flag if __name__ == "__main__": files = glob.glob('/folder/folder/2013/*.csv') num_workers = max(cpu_count()-2, 1) threshold = 75 for file in files: # 先只读需要的列建索引,降低内存占用 df = pd.read_csv(file, dtype=str, lineterminator='\n', usecols=['ID', 'text']) texts = df['text'].fillna("").tolist() inv_idx = build_inverted_index(texts) # 分块构造任务 chunk_size = len(texts) // num_workers + 1 tasks = [] for i in range(0, len(texts), chunk_size): chunk = texts[i:i+chunk_size] tasks.append((chunk, i, texts, inv_idx, threshold, 2)) # 多进程计算 all_flags = [] with Pool(num_workers) as pool: for res in pool.imap(check_chunk, tasks): all_flags.extend(res) # 最后读全表输出结果 df_full = pd.read_csv(file, dtype=str, lineterminator='\n') df_full[all_flags].to_csv(f'{file[:-4]}-done.csv', index=False) print(f'{file} 处理完成')
额外提效技巧
- 计算相似度前先对text做简单清洗:统一转小写、去掉多余空格、剔除推文中的固定链接、@用户名、特殊表情等无意义内容,既减少计算量,也能提升相似度判定的准确率
- 处理前先对text做全局去重,完全重复的文本只保留一份参与计算,避免重复劳动。
按上述方案优化后,单文件10万行级别的推文数据,在你现有的32核32G机器上,单文件处理时间可以从数天压缩到数分钟,不会再出现内存溢出问题。
内容的提问来源于stack exchange,提问作者Octner
相关产品推荐
相关产品推荐

