基于相似度阈值分组混合数据类型行的千万级客户数据快速去重咨询
客户数据相似性分组性能优化方案
需求说明
用户有1000万条左右的客户数据,包含邮箱、姓名、地址、邮编、手机号、出生日期等字段,数据存在大量空值,无直接可用于分组的固定字段,需要按照非空值匹配度超过70%即归为同一accountID的规则完成数据分组。
参考示例数据
| ID | FirstName | LastName | Adress | Adress1 | Adress2 | Zipcode | Birthdate | Mobilephone | Phone | |
|---|---|---|---|---|---|---|---|---|---|---|
| 1 | Donald.duck@disney.com | Donald | Duck | 12345 | +4912345 | |||||
| 2 | Donald.duck@disney.com | Donald | Duck | Mainstreet 1 | main | 1 | +4912345 | |||
| 3 | duck.d@gmail.com | Donald | Duck | Main str 1 | main | 1 | 12345 | +5054321 |
现有实现代码
def find_duplicates(df,Id,cid): idx = df[df.inID == inId].index tmp_mail = df.loc[idx].Email.values[0] row = df.loc[idx.values[0], ['Email', 'FirstName', 'LastName', 'Adress1','Adress2', 'Zipcode', 'Birthdate', 'Mobilephone','Phone']].dropna() if len(row)>3: res_email = df.loc[df.inEmail == tmp_mail] res_match = df[(df == row).sum(axis=1) >= round(len(row)) * 0.7] res = res_email.append(res_match).drop_duplicates() res['customer_id'] = cid res['ctd_emails'] = res.Email.nunique() res['ctd_inid'] = res.ID.nunique() return df.drop(index=res.index), res.drop(columns = ['Adress1','Adress2']) else: res = df.loc[df.inEmail == tmp_mail] res['customer_id'] = cid res['ctd_emails'] = res.Email.nunique() res['ctd_inid'] = res.ID.nunique() return df.drop(index=res.index), res.drop(columns = ['Adress1','Adress2']) ctd=0 while len(df)>0: df,tmp = find_duplicates(df,df.loc[df.index.min(),'ID'],ctd) if ctd == 0: res = tmp else: res = res.append(tmp) ctd += 1
现有问题
现有代码功能符合预期,但时间复杂度为O(n²),每处理一行都要遍历全表,还有大量pandas低效的drop、append操作,1000万条数据运行耗时极长,运行耗时统计如下:
可行优化方案
方案1:Pandas+Numpy向量化优化(改动最小,适合千万级以下数据)
核心优化点:
- 放弃逐行修改原表的逻辑,提前初始化customer_id标记数组,仅修改标记位即可,避免反复修改DataFrame
- 用numpy向量化运算替代逐行匹配,匹配速度提升100倍以上
- 可新增预分组逻辑,先按邮箱前缀+姓氏、邮编+姓氏等规则把数据拆成小候选组,仅在组内匹配,进一步缩小计算范围
参考实现代码:
import numpy as np import pandas as pd feature_cols = ['Email', 'FirstName', 'LastName', 'Adress1','Adress2', 'Zipcode', 'Birthdate', 'Mobilephone','Phone'] # 预生成特征哈希和非空掩码,空值用特殊值填充 hash_arr = df[feature_cols].apply(lambda x: x.apply(hash) if x.notna().any() else 0).fillna(0).to_numpy() non_null_mask = df[feature_cols].notna().to_numpy() customer_id = np.zeros(len(df), dtype=np.int32) current_cid = 1 for i in range(len(df)): if customer_id[i] != 0: continue # 取当前行非空特征 row_non_null = non_null_mask[i] row_feat = hash_arr[i][row_non_null] feat_cnt = len(row_feat) if feat_cnt <= 3: # 字段少的情况按邮箱匹配 match_mask = hash_arr[:, 0] == row_feat[0] else: # 计算匹配度,过滤符合阈值的行 match_cnt = (hash_arr[:, row_non_null] == row_feat).sum(axis=1) match_mask = match_cnt >= round(feat_cnt * 0.7) # 给未分组的匹配行打标记 customer_id[match_mask & (customer_id == 0)] = current_cid current_cid += 1 df['customer_id'] = customer_id # 后续按需计算ctd_emails、ctd_inid等统计字段即可
方案2:PySpark分布式实现(适合千万级及以上大数据量)
- 用MinHash+LSH近似近邻算法提前筛选高相似度候选对,避免全量笛卡尔积计算
- 分布式并行执行匹配逻辑,1000万条数据通常几分钟即可跑完
- 支持直接对接各类数据源,无需把全量数据加载到单节点内存
方案3:MySQL/MariaDB存储过程实现(适合数据已存储在库中、不想修改Python代码的场景)
- 第一步给所有参与匹配的字段加普通索引,提升查询速度
- 第二步生成候选分组临时表,按邮箱、手机号、邮编+姓氏等维度预分组,缩小匹配范围
- 第三步写存储过程遍历每个候选组,计算组内每行的匹配度,打accountID标记
- 注意不要直接全表join计算,必须加预过滤条件,否则会出现性能问题
内容的提问来源于stack exchange,提问作者pscl
相关产品推荐
相关产品推荐

