Python多线程优化百万级CSV相似度测试性能问题求助
优化大规模CSV数据相似度计算的性能问题
问题背景
处理116万行的CSV文件,单线程计算每行两两相似度耗时约7小时,尝试多线程优化但无效:
- 原相似度计算测试代码:
def similarity(): for i in range(0, 1000): for j in range(i+1, 1000): longestSentence = 0 commonWords = 0 row1 = dff['Product'].iloc[i] row2 = dff['Product'].iloc[j] wordsRow1 = row1.split() wordsRow2 = row2.split() # 找出两个句子中的相同单词 common = list(set(wordsRow1).intersection(wordsRow2)) if len(wordsRow1) > len(wordsRow2): longestSentence = len(wordsRow1) commonWords = calculate(common, wordsRow1) else: longestSentence = len(wordsRow2) commonWords = calculate(common, wordsRow2) print(i, j, (commonWords / longestSentence) * 100) def calculate(common, longestRow):#统计相同单词的总出现次数 sum = 0 for word in common: sum += longestRow.count(word) return sum
- 无效的多线程实现:
with ThreadPoolExecutor(max_workers=500) as executor: for result in executor.map(similarity()): print(result)
核心问题:设置大max_workers速度无变化,用threading库时多个线程重复执行整个函数。
核心优化方案
1. 修正多线程实现逻辑
你当前的代码直接执行了similarity(),并把执行结果传给executor.map,相当于线程没有拆分任务,而是重复执行整个函数。正确做法是把每一对(i,j)的计算拆成独立任务:
from concurrent.futures import ThreadPoolExecutor import pandas as pd # 预加载所有分词结果,避免反复调用iloc和split dff = pd.read_csv("your_file.csv") product_words = [row.split() for row in dff['Product']] def calculate_pair(i, j): wordsRow1 = product_words[i] wordsRow2 = product_words[j] common = set(wordsRow1).intersection(wordsRow2) if len(wordsRow1) > len(wordsRow2): longest_len = len(wordsRow1) common_count = sum(wordsRow1.count(word) for word in common) else: longest_len = len(wordsRow2) common_count = sum(wordsRow2.count(word) for word in common) return i, j, (common_count / longest_len) * 100 # 生成所有需要计算的(i,j)配对 pairs = [(i, j) for i in range(len(product_words)) for j in range(i+1, len(product_words))] # 线程数建议设为CPU核心数的2-4倍(Python线程受GIL限制,过大无意义) with ThreadPoolExecutor(max_workers=8) as executor: for result in executor.map(lambda x: calculate_pair(*x), pairs): # 建议写入文件而非打印,减少IO耗时 # with open("results.txt", "a") as f: # f.write(f"{result}\n") print(result)
2. 算法层面降维(关键优化)
原代码时间复杂度是O(n²*m)(n为行数,m为平均词数),116万行的n²是天文数字,这才是性能瓶颈,多线程无法解决:
- 预计算词频统计:用
Counter存储每个Product的词频,避免反复调用count():
from collections import Counter # 预计算每个Product的词频字典 product_word_counts = [Counter(row.split()) for row in dff['Product']] def calculate_pair_optimized(i, j): cnt1 = product_word_counts[i] cnt2 = product_word_counts[j] # 直接取共同单词在两个词频中的较小值求和,替代循环count common_count = sum(min(cnt1[word], cnt2[word]) for word in cnt1 if word in cnt2) longest_len = max(sum(cnt1.values()), sum(cnt2.values())) return i, j, (common_count / longest_len) * 100 if longest_len != 0 else 0.0
此优化将单对计算的时间复杂度从O(m)降到O(k)(k为共同单词数),大幅减少计算量。
- 避免全量两两比较:若只需找出相似度超过阈值的配对,用局部敏感哈希(LSH) 减少待计算配对数:
from datasketch import MinHash, MinHashLSH # 构建LSH索引,阈值设为你需要的相似度下限 lsh = MinHashLSH(threshold=0.5, num_perm=128) minhashes = [] for idx, words in enumerate(product_words): m = MinHash(num_perm=128) for word in words: m.update(word.encode('utf8')) minhashes.append(m) lsh.insert(idx, m) # 只查询可能相似的候选配对,避免全量n²计算 similar_pairs = set() for idx, m in enumerate(minhashes): candidates = lsh.query(m) for candidate in candidates: if candidate > idx: similar_pairs.add((idx, candidate)) # 仅对候选配对计算精确相似度 with ThreadPoolExecutor(max_workers=8) as executor: results = executor.map(lambda x: calculate_pair_optimized(*x), similar_pairs)
这能把待计算配对数从O(n²)降到O(n)级别,适合超大规模数据。
3. 改用多进程处理CPU密集任务
Python线程受GIL限制,CPU密集型任务无法真正并行,建议用ProcessPoolExecutor:
from concurrent.futures import ProcessPoolExecutor # 进程数建议设为CPU核心数 with ProcessPoolExecutor(max_workers=4) as executor: results = executor.map(lambda x: calculate_pair_optimized(*x), similar_pairs)
4. 细节优化
- 避免频繁
print:将结果写入文件而非打印,IO操作会严重拖慢速度。 - 分块处理:若内存不足,将CSV分成多个块,逐个处理后合并结果。
内容的提问来源于stack exchange,提问作者Berk Sunduri
相关产品推荐
相关产品推荐

