Python并行处理结果如何汇总至单个CSV?解决全局变量失效问题
高效并行处理CSV并汇总结果的方案
核心逻辑
别用全局变量,让每个并行任务直接返回处理后的结果数据,主进程收集所有返回值后统一合并成DataFrame再写入CSV——全程只做两次IO操作,比生成一堆小文件再合并高效太多,还能避开多进程内存隔离导致的全局变量失效问题。
方案1:用Python标准库concurrent.futures
无需额外安装,CPU密集型任务首选:
import pandas as pd from concurrent.futures import ProcessPoolExecutor def process_row(row): # 替换成你的实际处理逻辑,比如字段计算、数据转换 return { 'id': row['id'], 'processed_value': row['original_value'] * 2 # 示例操作 } if __name__ == '__main__': # 读取输入CSV input_df = pd.read_csv('input.csv') # 转成字典列表,方便逐行传递给并行任务 rows = input_df.to_dict('records') # 启动进程池处理 with ProcessPoolExecutor() as executor: # 用map分发任务,自动收集所有返回结果 results = list(executor.map(process_row, rows)) # 合并结果并写入最终CSV output_df = pd.DataFrame(results) output_df.to_csv('output.csv', index=False)
方案2:用joblib(如果你习惯这个库)
和上面逻辑一致,只是调用方式不同:
import pandas as pd from joblib import Parallel, delayed def process_row(row): # 替换为你的处理逻辑 return { 'id': row['id'], 'processed_value': row['original_value'] * 2 } if __name__ == '__main__': input_df = pd.read_csv('input.csv') rows = input_df.to_dict('records') # n_jobs=-1表示用满所有CPU核心 results = Parallel(n_jobs=-1)(delayed(process_row)(row) for row in rows) output_df = pd.DataFrame(results) output_df.to_csv('output.csv', index=False)
方案3:超大数据量优化(分块处理)
如果CSV大到内存装不下,就分块读+并行处理:
import pandas as pd from concurrent.futures import ProcessPoolExecutor def process_chunk(chunk): # 直接处理整个数据块,返回处理后的DataFrame chunk['processed_value'] = chunk['original_value'] * 2 # 示例操作 # 只保留需要的字段,减少内存占用 return chunk[['id', 'processed_value']] if __name__ == '__main__': # chunksize根据你的内存情况调整,比如1000行一块 chunk_iter = pd.read_csv('input.csv', chunksize=1000) with ProcessPoolExecutor() as executor: processed_chunks = list(executor.map(process_chunk, chunk_iter)) # 合并所有处理后的块 output_df = pd.concat(processed_chunks, ignore_index=True) output_df.to_csv('output.csv', index=False)
为什么你的全局变量方案没用?
Python多进程里,每个子进程会复制主进程的内存空间,子进程改全局变量只改自己副本里的,主进程根本看不到,所以最后output_df是空的。上面的方案都是让子进程把结果return出来,主进程主动收集,从根上解决了这个问题。
内容的提问来源于stack exchange,提问作者Ash
相关产品推荐
相关产品推荐

