基于批处理的多进程优化: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
相关产品推荐
相关产品推荐

