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

大数据文件模糊匹配多进程运行异常问题求助

多进程模糊匹配陷入无限运行问题排查

我有两个数据文件:一个含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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 18:55:57