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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 07:45:58