如何加速大数据集下两个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
相关产品推荐
相关产品推荐

