Python多进程并行处理多CSV文件的高效方案及最优分片数咨询
高效处理大批量CSV文件的Python多进程方案
这问题我太熟了——手动分片存pkl、开多个终端确实够折腾的,尤其是文件多负载重的时候,完全可以用Python原生工具把整个流程自动化,不用再手动搞这些繁琐操作!下面给你详细说下可行方案和最优分片数的确定方法:
核心优化思路
你当前的手动分片方式本质是手动实现并行,但没法自动调度系统资源。Python自带的多进程库可以帮你自动完成任务拆分、进程管理、结果合并这整套流程,完全替代手动操作。
推荐的并行处理方案
1. concurrent.futures.ProcessPoolExecutor(最省心的选择)
这个库是Python3.2+原生自带的,封装了多进程的底层细节,不用自己管进程启动/销毁,直接提交任务就行。特别适合pandas处理CSV这类CPU密集型任务。
2. 原生multiprocessing模块(更灵活的控制)
如果需要更精细的进程管理(比如进程间通信、自定义进程初始化),可以用这个模块,不过封装性不如前者,代码量会稍多一点。
3. Dask(超大数据集场景)
如果你的CSV总数据量已经大到单进程内存扛不住,Dask是更好的选择——它完美兼容pandas API,自动分片并行处理,还能处理内存外的超大文件。
如何确定最优分片数?
分片数(也就是进程数)的选择要结合任务类型和硬件资源:
- CPU密集型任务:最优进程数通常等于你的CPU物理核心数(或者核心数+1,用来利用CPU空闲等待时间)。可以用
os.cpu_count()直接获取核心数。 - IO密集型任务:如果处理中大部分时间在读取文件(磁盘IO等待),可以把进程数设为核心数的2-4倍,这样某个进程等待IO时,其他进程能继续工作。
- 实际测试校准:最好自己测试不同进程数的处理耗时,比如从核心数的1倍到3倍,看哪个速度最快——有时候受限于内存或磁盘IO,不是进程越多效率越高。
代码示例
用ProcessPoolExecutor实现自动并行处理
import os import numpy as np import pandas as pd from concurrent.futures import ProcessPoolExecutor # 定义单个CSV文件的处理逻辑 def process_single_file(file_path): # 这里替换成你的实际处理代码 df = pd.read_csv(file_path) df['cleaned_col'] = df['raw_col'].str.strip() df['calculated_col'] = df['num_col'] * 1.5 return df # 定义分片处理函数:处理一个文件列表,合并为单个DataFrame def process_chunk(file_list): dfs = [process_single_file(f) for f in file_list] return pd.concat(dfs, ignore_index=True) if __name__ == '__main__': # 你的原始文件列表 master_list = ["file1.csv", "file2.csv", ..., "filen.csv"] # 获取CPU核心数作为进程数 num_processes = os.cpu_count() # 自动拆分文件列表为对应分片 all_chunks = np.array_split(master_list, num_processes) # 启动进程池并行处理 with ProcessPoolExecutor(max_workers=num_processes) as executor: # 提交所有分片任务,获取处理结果 chunk_results = list(executor.map(process_chunk, all_chunks)) # 合并所有分片结果为最终DataFrame final_df = pd.concat(chunk_results, ignore_index=True) # 保存最终结果 final_df.to_csv("merged_final.csv", index=False) # 也可以保存为pkl # final_df.to_pickle("merged_final.pkl")
关键注意事项
- 一定要把主逻辑放在
if __name__ == '__main__':代码块下,这是Windows系统多进程的强制要求,Linux/macOS虽然不强制,但加上能避免很多潜在问题。 - 如果处理函数需要额外参数,可以用
functools.partial包装后再提交。 - 如果内存紧张,可以让每个进程处理完一个文件就保存临时结果,最后再合并临时文件,避免把所有DataFrame都放在内存里。
替代方案:用multiprocessing.Pool
和ProcessPoolExecutor逻辑类似,写法稍有不同:
import multiprocessing as mp import numpy as np import pandas as pd def process_chunk(file_list): dfs = [] for f in file_list: df = pd.read_csv(f) # 你的处理逻辑... dfs.append(df) return pd.concat(dfs, ignore_index=True) if __name__ == '__main__': master_list = ["file1.csv", ...] num_processes = mp.cpu_count() chunks = np.array_split(master_list, num_processes) with mp.Pool(num_processes) as pool: results = pool.map(process_chunk, chunks) final_df = pd.concat(results, ignore_index=True) final_df.to_csv("merged_result.csv")
Dask处理超大数据集示例
如果数据量已经超出单进程内存,用Dask几乎不用改pandas代码:
import dask.dataframe as dd # 直接读取所有CSV,Dask自动分片并行处理 ddf = dd.read_csv("*.csv") # 执行你的处理逻辑,语法和pandas完全一致 ddf['cleaned_col'] = ddf['raw_col'].str.strip() # 计算并保存结果(如果内存够可以用compute()转成pandas DataFrame,否则直接保存) ddf.to_csv("dask_merged_*.csv", index=False)
这样你就不用再手动分片、存pkl、开终端了,整个流程自动并行处理,效率提升不止一个档次!
内容的提问来源于stack exchange,提问作者Luis Miguel
相关产品推荐
相关产品推荐

