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

基于批处理的多进程优化:3000万行大文件内存问题求解

多进程处理大文件(3000万行)+ 避免内存溢出方案

针对你的需求,核心解决思路是分块读取+并行处理+安全写入,既利用多进程加速,又从根源避免全量加载数据导致的内存溢出。以下是具体实现和优化点:

核心问题修正

你的现有代码已经用到分批次处理,但前提是df2_csv是全量加载的3000万行DataFrame,这才是内存溢出的主要诱因。必须先改成分块读取原文件,而非一次性加载全量数据。

方案一:进程池+统一写入(简单高效)

import pandas as pd
from multiprocessing import Pool
import tqdm

# 假设你的match_and_merge业务逻辑已定义
def match_and_merge(batch):
    # 你的数据匹配/合并逻辑
    pass

# 单个进程处理逻辑:接收csv分块,返回处理后的DataFrame
def process_chunk(chunk):
    batch = chunk.to_dict(orient='records')
    results = match_and_merge(batch)
    return pd.DataFrame(results, columns=matched_df.columns)

if __name__ == '__main__':
    INPUT_FILE = "你的输入文件路径.csv"  # 替换为实际输入文件
    OUTPUT_FILE = "final_patent_sample.csv"
    CHUNK_SIZE = 10000  # 可根据内存调整,推荐1-5万行/块
    PROCESS_NUM = 10  # 你计划的10个进程

    # 1. 分块读取原文件,不加载全量数据
    csv_chunks = pd.read_csv(INPUT_FILE, chunksize=CHUNK_SIZE)
    total_chunks = 30000000 // CHUNK_SIZE  # 预估总块数,用于进度条

    # 2. 进程池并行处理
    with Pool(processes=PROCESS_NUM) as pool:
        # imap比map更省内存,逐个提交任务
        processed_results = list(tqdm.tqdm(pool.imap(process_chunk, csv_chunks), total=total_chunks))

    # 3. 统一写入结果文件
    header_written = False
    for result_df in processed_results:
        result_df.to_csv(OUTPUT_FILE, mode='a', index=False, header=not header_written)
        header_written = True

方案二:临时文件+合并(更安全,适合超大规模数据)

如果担心中间进程出错导致数据丢失,或者单个处理后的DataFrame依然占用过多内存,可以让每个进程写入临时文件,最后统一合并:

import pandas as pd
from multiprocessing import Pool
import tqdm
import os

def match_and_merge(batch):
    # 你的业务逻辑
    pass

# 单个进程处理并写入临时文件
def process_chunk_to_temp(chunk_idx, chunk):
    batch = chunk.to_dict(orient='records')
    results = match_and_merge(batch)
    temp_file_path = f"./temp_chunks/chunk_{chunk_idx}.csv"
    pd.DataFrame(results, columns=matched_df.columns).to_csv(temp_file_path, index=False, header=False)
    return temp_file_path

if __name__ == '__main__':
    INPUT_FILE = "你的输入文件路径.csv"
    OUTPUT_FILE = "final_patent_sample.csv"
    CHUNK_SIZE = 10000
    PROCESS_NUM = 10
    TEMP_DIR = "./temp_chunks"

    # 创建临时目录
    os.makedirs(TEMP_DIR, exist_ok=True)

    # 分块读取并带上索引,方便生成唯一临时文件名
    csv_chunks = enumerate(pd.read_csv(INPUT_FILE, chunksize=CHUNK_SIZE))
    total_chunks = 30000000 // CHUNK_SIZE

    # 并行处理生成临时文件
    with Pool(processes=PROCESS_NUM) as pool:
        temp_files = list(tqdm.tqdm(pool.starmap(process_chunk_to_temp, csv_chunks), total=total_chunks))

    # 合并临时文件到最终输出
    with open(OUTPUT_FILE, 'w', encoding='utf-8') as out_file:
        # 先写入表头
        out_file.write(','.join(matched_df.columns) + '\n')
        # 逐个写入临时文件内容
        for temp_file in temp_files:
            with open(temp_file, 'r', encoding='utf-8') as tf:
                out_file.write(tf.read())
            os.remove(temp_file)
    
    # 删除临时目录
    os.rmdir(TEMP_DIR)

关键优化点说明

  • 分块读取:用pd.read_csv(chunksize=...)替代全量加载,每次仅将一小块数据载入内存,从根源避免内存溢出。
  • 进程池选择:使用Pool.imap或starmap,避免一次性将所有分块放入任务队列,进一步降低内存占用。
  • 写入安全:方案一用单进程统一追加写入,避免多进程同时写文件导致的内容错乱;方案二用临时文件,即使部分进程失败,也不会丢失全部已处理数据。
  • 参数调整:根据你的内存大小调整CHUNK_SIZE——内存充足可以调大(如5万行),减少进程调度开销;内存紧张则调小。

注意事项

  • 确保match_and_merge函数可序列化:不能包含无法被pickle序列化的对象(如全局的数据库连接、打开的文件句柄),这类资源需要在函数内部重新初始化。
  • 测试先行:先用小样本数据验证逻辑正确性,再运行全量数据,避免浪费时间。

内容的提问来源于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 17:54:58