Python生成大型组合DataFrame时的内存问题解决方案咨询
解决8万条词汇模糊相似度匹配的内存溢出问题
核心问题
8万条去重词汇的两两组合约32亿行,直接生成全量DataFrame会耗尽内存(即使32GB RAM也无法承载),必须采用流式处理、分块计算、预过滤等策略避免内存堆积。
方案1:流式分块处理,直接写入磁盘
不一次性生成全量组合DataFrame,而是迭代组合对,每处理一定数量的组合就写入磁盘,实时释放内存。
import itertools import pandas as pd from thefuzz import fuzz def process_and_write_chunk(chunk, output_path, mode='a', header=False): """处理单块组合并写入文件""" df_chunk = pd.DataFrame(chunk, columns=['name1', 'name2']) # 计算模糊评分 df_chunk[['partial_ratio', 'ratio']] = df_chunk.apply( lambda row: pd.Series([ fuzz.partial_ratio(row['name1'], row['name2']), fuzz.ratio(row['name1'], row['name2']) ]), axis=1 ) df_chunk.to_csv(output_path, mode=mode, header=header, index=False) # 配置参数 unique_words = [你的8万条去重词汇列表] chunk_size = 1_000_000 # 每100万条组合为一个块,可根据内存调整 output_file = 'fuzzy_matches.csv' # 先写入表头 pd.DataFrame(columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='w', header=True, index=False ) # 流式迭代组合并处理 current_chunk = [] for idx, (word1, word2) in enumerate(itertools.combinations(unique_words, 2)): current_chunk.append((word1, word2)) # 达到块大小则处理写入 if (idx + 1) % chunk_size == 0: process_and_write_chunk(current_chunk, output_file) current_chunk = [] # 清空块释放内存 # 处理剩余的最后一批数据 if current_chunk: process_and_write_chunk(current_chunk, output_file)
方案2:用Dask进行分布式并行计算
Dask能自动将大数据集分块,并行处理并写入磁盘,无需手动管理分块,适合大规模数据任务。
import dask.dataframe as dd import itertools from thefuzz import fuzz unique_words = [你的8万条去重词汇列表] # 生成组合迭代器,Dask可直接处理 combinations_iter = itertools.combinations(unique_words, 2) # 转为Dask DataFrame,根据CPU核心数设置分区数(i7-1185G7建议设为16) ddf = dd.from_pandas( pd.DataFrame(combinations_iter, columns=['name1', 'name2']), npartitions=16 ) # 定义模糊评分计算函数 def calculate_fuzzy_scores(row): return pd.Series([ fuzz.partial_ratio(row['name1'], row['name2']), fuzz.ratio(row['name1'], row['name2']) ], index=['partial_ratio', 'ratio']) # 应用函数并计算(Dask自动并行处理) ddf = ddf.assign(**ddf.apply(calculate_fuzzy_scores, axis=1, meta={ 'partial_ratio': int, 'ratio': int })) # 写入Parquet格式(比CSV更高效,支持后续快速查询) ddf.to_parquet('fuzzy_matches.parquet', write_index=False)
方案3:预过滤减少不必要计算
大部分词汇之间相似度极低,提前通过长度差、N-gram交集筛选出可能相似的词汇,大幅减少计算量。
import itertools import pandas as pd from thefuzz import fuzz from collections import defaultdict def get_2grams(word): """生成词汇的2-gram集合""" if len(word) < 2: return set() return set([word[i:i+2] for i in range(len(word)-1)]) unique_words = [你的8万条去重词汇列表] output_file = 'fuzzy_matches_filtered.csv' # 按词汇长度分组,只比较长度差≤2的词汇 length_groups = defaultdict(list) for word in unique_words: length_groups[len(word)].append(word) # 写入表头 pd.DataFrame(columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='w', header=True, index=False ) # 处理同长度组内的组合 for length in length_groups: group = length_groups[length] for word1, word2 in itertools.combinations(group, 2): # 2-gram交集占比≥0.5才计算模糊评分(可调整阈值) grams1 = get_2grams(word1) grams2 = get_2grams(word2) if not grams1 or not grams2: continue overlap_ratio = len(grams1 & grams2) / min(len(grams1), len(grams2)) if overlap_ratio >= 0.5: pr = fuzz.partial_ratio(word1, word2) r = fuzz.ratio(word1, word2) pd.DataFrame([[word1, word2, pr, r]], columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='a', header=False, index=False ) # 处理长度差为1、2的跨组组合 for length in length_groups: for delta in [1, 2]: target_length = length + delta if target_length not in length_groups: continue group1 = length_groups[length] group2 = length_groups[target_length] for word1 in group1: for word2 in group2: grams1 = get_2grams(word1) grams2 = get_2grams(word2) if not grams1 or not grams2: continue overlap_ratio = len(grams1 & grams2) / min(len(grams1), len(grams2)) if overlap_ratio >= 0.5: pr = fuzz.partial_ratio(word1, word2) r = fuzz.ratio(word1, word2) pd.DataFrame([[word1, word2, pr, r]], columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='a', header=False, index=False )
方案4:多进程队列式处理
用多进程拆分计算任务,主进程生成组合放入队列,子进程计算后写入结果,避免主进程内存堆积。
import itertools import pandas as pd from thefuzz import fuzz from multiprocessing import Pool, Manager def worker_task(pair_queue, result_queue): """子进程处理组合,计算评分后放入结果队列""" while True: pair = pair_queue.get() if pair is None: # 结束信号 break word1, word2 = pair pr = fuzz.partial_ratio(word1, word2) r = fuzz.ratio(word1, word2) result_queue.put((word1, word2, pr, r)) def writer_task(result_queue, output_file): """写入进程,从结果队列取数据写入磁盘""" pd.DataFrame(columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='w', header=True, index=False ) while True: result = result_queue.get() if result is None: # 结束信号 break pd.DataFrame([result], columns=['name1', 'name2', 'partial_ratio', 'ratio']).to_csv( output_file, mode='a', header=False, index=False ) if __name__ == '__main__': unique_words = [你的8万条去重词汇列表] output_file = 'fuzzy_matches_multiprocess.csv' num_workers = 4 # 根据CPU核心数调整(i7-1185G7建议4-8) with Manager() as manager: pair_queue = manager.Queue(maxsize=1000) # 限制队列大小避免内存溢出 result_queue = manager.Queue(maxsize=1000) # 启动写入进程 writer_proc = manager.Process(target=writer_task, args=(result_queue, output_file)) writer_proc.start() # 启动工作进程池 with Pool(num_workers, worker_task, (pair_queue, result_queue)) as pool: # 生成组合并放入队列 for pair in itertools.combinations(unique_words, 2): pair_queue.put(pair) # 发送结束信号给所有工作进程 for _ in range(num_workers): pair_queue.put(None) # 发送结束信号给写入进程 result_queue.put(None) writer_proc.join()
内容的提问来源于stack exchange,提问作者Rohan
相关产品推荐
相关产品推荐

