大数据文件模糊匹配多进程运行异常问题求助
多进程模糊匹配陷入无限运行问题排查
我有两个数据文件:一个含100万行数据(30MB),另一个含3000万行数据(9GB)。需要基于companyname_std列和disambig_assignee_organization_std列进行模糊匹配。单进程运行耗时极长,尝试用多进程加速后,即便处理样本文件也会一直运行,不得不终止进程,求排查原因。
单进程代码
import pandas as pd from rapidfuzz import process from tqdm import tqdm import re import multiprocessing as mp # 标准化组织名称 def standardize_text(text): if pd.isnull(text): return '' # 转小写 text = text.lower() # 移除特殊字符 text = re.sub(r'[^\w\s]', '', text) # 移除多余空格 text = ' '.join(text.split()) return text def match_and_merge(rows): results = [] for row in rows: match = process.extract(row['companyname_std'], assignee_orgs_std, limit=1) if match and match[0][1] > 90: matched_org = df1_csv[df1_csv['disambig_assignee_organization_std'] == match[0][0]].iloc[0] row_series = pd.Series(row) merged_row = pd.concat([row_series, matched_org], axis=0) results.append(merged_row.tolist() + [1]) else: row_series = pd.Series(row) empty_series = pd.Series([None] * len(df1_csv.columns), index=df1_csv.columns) merged_row = pd.concat([row_series, empty_series], axis=0) results.append(merged_row.tolist() + [0]) return results # 文件路径 csv_file1 = "g_patent.csv" csv_file2 = "1. pitchbook_company.csv" output_file = "final_patent_sample.csv" # 加载数据(样本取前1000行) df1_csv = pd.read_csv(csv_file1, nrows=1000) df2_csv = pd.read_csv(csv_file2, nrows=1000) # 标准化名称列 df1_csv['disambig_assignee_organization_std'] = df1_csv['disambig_assignee_organization'].apply(standardize_text) df2_csv['companyname_std'] = df2_csv['companyname'].apply(standardize_text) # 提取用于匹配的标准化列表 assignee_orgs_std = df1_csv['disambig_assignee_organization_std'].tolist() # 初始化结果文件或续传 try: matched_df = pd.read_csv(output_file) start_index = matched_df.shape[0] print(f"从索引 {start_index} 开始续传") except FileNotFoundError: matched_df = pd.DataFrame(columns=df2_csv.columns.tolist() + df1_csv.columns.tolist() + ['merge_flag']) start_index = 0 print("未找到历史文件,从头开始处理") # 截断数据到续传位置 df2_csv = df2_csv[start_index:] batch_size = 100 header_written = False # 单进程批量处理 for i in tqdm(range(0, df2_csv.shape[0], batch_size)): batch = df2_csv.iloc[i:i + batch_size].to_dict(orient='records') results = match_and_merge(batch) batch_df = pd.DataFrame(results, columns=matched_df.columns) batch_df.to_csv(output_file, mode='a', index=False, header=not header_written) header_written = True
多进程代码片段
with multiprocessing.Pool(10) as pool: batch_generator = (df2_csv.iloc[i:i + batch_size].to_dict(orient='records') for i in range(0, df2_csv.shape[0], batch_size)) for results in pool.imap(match_and_merge, batch_generator): batch_df = pd.DataFrame(results, columns=matched_df.columns) batch_df.to_csv(output_file, mode='a', index=False, header=not header_written) header_written = True
问题根源
- 全局变量跨进程冗余复制:
match_and_merge依赖全局变量assignee_orgs_std和df1_csv,每个子进程启动时会完整复制这些数据,样本量小时开销不明显,但实际数据量大会导致内存过载、进程初始化耗时极长,表现为“一直运行”。 - 进程间通信开销过大:每次传递的批次是完整的字典列表,加上全局数据的复制,进程间数据传输耗时远超计算时间,导致任务迟迟无法完成。
- 文件写入竞态:多进程同时追加写入同一个CSV文件,底层文件锁可能导致进程阻塞,看似无限运行。
- 逐行匹配效率低下:即便用了多进程,仍采用逐行
process.extract,未利用rapidfuzz的批量优化接口,计算效率未达标。
修复方案
1. 重构函数,避免全局变量传递
提前构建匹配映射,将必要数据作为参数传入子进程,减少冗余复制:
# 主进程中提前构建匹配映射,避免子进程重复查询DataFrame org_map = df1_csv.set_index('disambig_assignee_organization_std').to_dict('index') def match_and_merge(args): rows, org_list, org_map = args results = [] for row in rows: match = process.extract(row['companyname_std'], org_list, limit=1) if match and match[0][1] > 90: matched_data = org_map.get(match[0][0], {}) merged_row = {**row, **matched_data, 'merge_flag': 1} else: merged_row = {**row, **{col: None for col in df1_csv.columns}, 'merge_flag': 0} results.append(merged_row) return results # 进程池调用时传递参数 with mp.Pool(mp.cpu_count()-1) as pool: task_generator = ( (df2_csv.iloc[i:i+batch_size].to_dict('records'), assignee_orgs_std, org_map) for i in range(0, df2_csv.shape[0], batch_size) ) header_written = False for results in pool.imap(match_and_merge, task_generator): batch_df = pd.DataFrame(results, columns=matched_df.columns) # 加锁保证写入安全 with open(output_file, 'a', newline='', encoding='utf-8') as f: batch_df.to_csv(f, index=False, header=not header_written) header_written = True
2. 改用批量模糊匹配优化计算
使用process.cdist批量处理整个批次,替代逐行循环:
def match_and_merge_batch(args): batch_df, org_list, org_map = args # 批量计算相似度,直接过滤得分≥90的结果 scores = process.cdist(batch_df['companyname_std'], org_list, score_cutoff=90) results = [] for idx, row in batch_df.iterrows(): matches = scores[idx] if matches: # 取最高分匹配项 best_match = max(matches, key=lambda x: x[1]) matched_data = org_map.get(best_match[0], {}) merged_row = {**row.to_dict(), **matched_data, 'merge_flag': 1} else: merged_row = {**row.to_dict(), **{col: None for col in df1_csv.columns}, 'merge_flag': 0} results.append(merged_row) return results # 传递批次DataFrame而非字典列表,减少序列化开销 task_generator = ( (df2_csv.iloc[i:i+batch_size], assignee_orgs_std, org_map) for i in range(0, df2_csv.shape[0], batch_size) )
3. 合理设置进程数
避免硬编码进程数,改用CPU核心数减1,防止资源过载:
with mp.Pool(mp.cpu_count()-1) as pool: # 后续代码
4. 优化大文件加载
对于9GB的大文件,不要一次性加载到内存,改用pandas.read_csv的chunksize参数分块加载,进一步减少内存压力。
内容的提问来源于stack exchange,提问作者Cindy
相关产品推荐
相关产品推荐

