基于RapidFuzz的海量新闻去重:2万+记录性能优化问询
问题
我正在使用RapidFuzz的token_set_ratio检测Refinitiv数据中的近重复新闻文章。数据集规模可超过2万条记录,目前采用多进程方式对比记录对。该逻辑在小输入(如5千条)下运行良好,但在大输入时系统会挂起或无响应。
我将所有两两组合拆分为块,通过ProcessPoolExecutor并行处理,使用如下函数:
def process_chunk(chunk): results = [] for rec1, rec2 in chunk: text1 = f"{rec1['data'].get('headline', '')}. {rec1['data'].get('body', '')}" text2 = f"{rec2['data'].get('headline', '')}. {rec2['data'].get('body', '')}" score = token_set_ratio(text1, text2) if score > threshold: results.append((rec1['guid'], rec2['guid'], score)) return results
恳请提供以下方面的建议:
- 如何高效扩展至数万条记录?
- 如何减少两两对比时的内存/CPU瓶颈?
- 能否用更智能的预过滤替代全量两两对比?
解决方案建议
1. 高效扩展至数万条记录的方案
- 优化进程池配置:设置
max_workers为CPU核心数的1-2倍,避免进程切换开销过大;同时调整chunksize,按每条进程处理100-500对记录的粒度拆分任务,减少进程通信的额外开销。 - 分批处理并即时落地结果:不要等待所有任务完成再收集结果,用
as_completed逐个获取chunk的计算结果,立即写入文件或数据库,避免大量结果堆积在内存导致OOM。示例逻辑:from concurrent.futures import ProcessPoolExecutor, as_completed with ProcessPoolExecutor(max_workers=8) as executor: futures = [executor.submit(process_chunk, chunk) for chunk in chunk_list] for future in as_completed(futures): partial_results = future.result() # 即时写入结果,释放内存 with open('duplicate_results.csv', 'a') as f: for guid1, guid2, score in partial_results: f.write(f"{guid1},{guid2},{score}\n") - 避免重复计算:只生成
i < j的记录对(即不重复对比rec1-rec2和rec2-rec1),直接减少一半的计算量。
2. 减少内存/CPU瓶颈的方法
- 预处理文本并缓存:提前拼接每条记录的
headline+body并缓存,避免在每个配对中重复执行字符串拼接操作。比如:
后续在text_cache = { rec['guid']: f"{rec['data'].get('headline', '')}. {rec['data'].get('body', '')}" for rec in all_records }process_chunk中直接通过guid从缓存取文本即可。 - 简化文本复杂度:对文本做预处理,比如转小写、去除标点/停用词、提取核心关键词,减少
token_set_ratio的计算量。示例:import re from nltk.corpus import stopwords stop_words = set(stopwords.words('english')) def preprocess(text): # 去标点、转小写 text = re.sub(r'[^\w\s]', '', text.lower()) # 去停用词 return ' '.join([w for w in text.split() if w not in stop_words]) # 预处理并缓存所有文本 text_cache = {rec['guid']: preprocess(f"{rec['data'].get('headline', '')}. {rec['data'].get('body', '')}") for rec in all_records} - 改用线程池(可选):如果RapidFuzz的
token_set_ratio是C实现(主流版本均为),GIL不会成为瓶颈,改用ThreadPoolExecutor可减少进程间数据复制的内存开销。 - 收紧相似度阈值:设置更高的
threshold,过滤掉低相似度的记录对,减少结果数据量和无效计算。
3. 替代全量两两对比的预过滤方案
- MinHash/LSH分簇:用MinHash将文本转换为哈希签名,再通过LSH将相似文本分到同一簇,只对比同簇内的记录对,可大幅减少配对数量。可以用
datasketch库快速实现。 - 元数据过滤:利用Refinitiv新闻的元数据提前过滤,比如只对比发布时间差在24小时内、同一来源或同一分类的新闻,直接排除明显不可能重复的配对。
- 关键词交集过滤:提取每条新闻的TopN高频关键词,只有当两条新闻的关键词交集达到一定数量(比如≥3个)时,才执行
token_set_ratio计算,用集合交集判断的速度远快于相似度计算。 - TF-IDF+KMeans分簇:先计算所有文本的TF-IDF向量,用KMeans将文本分成若干主题簇,只对比同簇内的记录,适合主题明确的新闻数据。
内容的提问来源于stack exchange,提问作者Werner Spreeuwenberg
相关产品推荐
相关产品推荐

