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

如何加速大数据集下两个DataFrame的地址匹配及置信度计算?

优化大规模地址模糊匹配的方法

针对你300万条主数据+10条参考数据的地址匹配场景,以下是针对性的提速方案,直接解决核心性能问题:


一、砍掉冗余计算:用extractOne替代extract

你的代码里process.extract默认返回Top5匹配结果,但你只需要最高分匹配,完全没必要计算多余结果。替换成process.extOne能直接减少80%的计算量:

# 替换原match_addresses函数
def match_addresses(add, contacts_addresses, min_score=0):
    # 直接返回最高分匹配的元组:(匹配值, 分数, 索引)
    return process.extOne(add, contacts_addresses, scorer=fuzz.token_sort_ratio, score_cutoff=min_score)

同时可以直接删掉get_highest_score函数,因为extractOne已经给出了最优结果。


二、真正利用Dask的并行能力

你用了Dask加载主数据,但又转成list(df.address)把所有数据塞进单进程内存,完全浪费了Dask的分块并行优势。改成基于Dask分区的并行计算:

import pandas as pd
from rapidfuzz import process, fuzz
from dask import dataframe as dd

# 加载参考数据
ref_df = pd.read_csv('reference_dataset.csv')
ref_addresses = ref_df.ref_address.unique().tolist()

# 用Dask分块加载主数据(默认按文件大小分块,自动利用多进程)
df = dd.read_csv('main_dataset.csv', low_memory=False)

def find_best_in_partition(partition, ref_addrs):
    results = []
    for ref_addr in ref_addrs:
        best_match = process.extOne(ref_addr, partition['address'], scorer=fuzz.token_sort_ratio, score_cutoff=75)
        if best_match:
            results.append({
                'ref_address': ref_addr,
                'matched_address': best_match[0],
                'score': int(best_match[1])
            })
    return pd.DataFrame(results)

# 对每个数据分区并行计算匹配结果
partition_results = df.map_partitions(find_best_in_partition, ref_addresses, meta={
    'ref_address': str,
    'matched_address': str,
    'score': int
})

# 合并所有分区结果,取每个参考地址的最高分匹配
final_results = partition_results.groupby('ref_address').apply(
    lambda x: x.loc[x['score'].idxmax()],
    meta={
        'ref_address': str,
        'matched_address': str,
        'score': int
    }
).compute()

# 和原参考表合并输出
merged_results = pd.merge(ref_df, final_results, how='right', on='ref_address')
merged_results.to_csv('results.csv', index=False)

这种方式不需要把300万条数据全部加载到内存,Dask会自动分配到多个进程并行处理,速度能提升数倍。


三、地址预处理:减少模糊匹配的计算量

模糊匹配的速度和字符串复杂度直接相关,先对地址做标准化处理,既能提速又能提升匹配准确度:

def clean_address(address):
    if pd.isna(address):
        return ''
    # 转小写、去特殊符号
    addr = str(address).lower().replace(',', '').replace('.', '').replace('-', '')
    # 标准化常见缩写
    addr = addr.replace(' st ', ' street ').replace(' ave ', ' avenue ')
    addr = addr.replace(' rd ', ' road ').replace(' blvd ', ' boulevard ')
    # 去除多余空格
    return ' '.join(addr.split())

# 对主数据和参考数据的地址做清洗
df['clean_address'] = df['address'].apply(clean_address, meta=str)
ref_df['clean_ref_address'] = ref_df['ref_address'].apply(clean_address)
ref_addresses_clean = ref_df['clean_ref_address'].unique().tolist()

# 后续匹配用清洗后的地址(记得存回原始地址)
def find_best_in_partition(partition, ref_addrs_clean):
    results = []
    for clean_ref in ref_addrs_clean:
        best_match = process.extOne(clean_ref, partition['clean_address'], scorer=fuzz.token_sort_ratio, score_cutoff=75)
        if best_match:
            # 找回原始参考地址和匹配地址
            original_ref = ref_df.loc[ref_df['clean_ref_address'] == clean_ref, 'ref_address'].iloc[0]
            original_match = partition.iloc[best_match[2]]['address']
            results.append({
                'ref_address': original_ref,
                'matched_address': original_match,
                'score': int(best_match[1])
            })
    return pd.DataFrame(results)

清洗后的字符串更短、更标准,模糊匹配的计算量会大幅降低。


四、批量匹配:用cdist替代循环

RapidFuzz的process.cdist可以批量计算两个字符串列表的相似度矩阵,结合Dask分块能进一步提速:

from rapidfuzz import process, fuzz

def compute_similarity(partition, ref_addrs_clean):
    partition_clean = partition['clean_address'].tolist()
    # 批量计算相似度矩阵,自动用多进程加速
    matrix = process.cdist(ref_addrs_clean, partition_clean, scorer=fuzz.token_sort_ratio, workers=-1)
    results = []
    for i, clean_ref in enumerate(ref_addrs_clean):
        max_score = matrix[i].max()
        max_idx = matrix[i].argmax()
        original_ref = ref_df.loc[ref_df['clean_ref_address'] == clean_ref, 'ref_address'].iloc[0]
        original_match = partition.iloc[max_idx]['address']
        results.append({
            'ref_address': original_ref,
            'matched_address': original_match,
            'score': int(max_score)
        })
    return pd.DataFrame(results)

# 应用到每个分区
partition_results = df.map_partitions(compute_similarity, ref_addresses_clean, meta={
    'ref_address': str,
    'matched_address': str,
    'score': int
})

cdist的批量计算比循环单条处理效率更高,适合你的参考数据量小、主数据量大的场景。


其他小优化

  • 不要把Dask DataFrame转成list:list(df.address)会把所有数据加载到内存,300万条数据可能触发内存交换导致变慢;
  • 合理设置score_cutoff:比如设为75,直接跳过低分数匹配,减少无效计算;
  • 更换更快的scorer:如果token_sort_ratio速度不够,可尝试ratio或partial_ratio(注意平衡准确度)。

内容的提问来源于stack exchange,提问作者Kelly Tang

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:15:33