You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

恳请提供以下方面的建议:

  1. 如何高效扩展至数万条记录?
  2. 如何减少两两对比时的内存/CPU瓶颈?
  3. 能否用更智能的预过滤替代全量两两对比?

解决方案建议

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.13 07:42:09