如何利用多数据块并行执行相同处理流程?
并行处理DataFrame分块的实现方案
首先你需要先定义单个数据块的处理函数,比如针对每个块里的文件名执行读取、分析等操作,示例如下:
import pandas as pd import os def process_chunk(chunk_df, wb_path): # 这里写你对单个数据块的处理逻辑,比如读取每个文件并处理 results = [] for filename in chunk_df['Wk_nm']: file_path = os.path.join(wb_path, filename) # 示例:读取Excel文件(根据你的实际文件类型调整) df_file = pd.read_excel(file_path) # 这里添加你的处理逻辑,比如计算统计值、数据清洗等 chunk_result = { 'filename': filename, 'row_count': len(df_file), 'mean_value': df_file['某列'].mean() if '某列' in df_file.columns else None } results.append(chunk_result) # 返回该块的处理结果 return pd.DataFrame(results)
接下来提供几种常用的并行实现方式:
方法1:使用concurrent.futures.ProcessPoolExecutor(推荐,语法简洁)
from concurrent.futures import ProcessPoolExecutor if __name__ == '__main__': wb_path = "你的目录路径" # 你的分块代码(已给出) entries = os.listdir(wb_path) df = pd.DataFrame(entries, columns=['Wk_nm']) n = 50 list_df = [df[i:i+n] for i in range(0, len(df), n)] # 并行处理 with ProcessPoolExecutor() as executor: # 提交所有分块的处理任务,每个任务传入分块和目录路径 futures = [executor.submit(process_chunk, chunk, wb_path) for chunk in list_df] # 收集所有结果并合并 all_results = pd.concat([future.result() for future in futures], ignore_index=True) # 查看最终合并后的结果 print(all_results)
方法2:使用multiprocessing.Pool(传统多进程方式)
from multiprocessing import Pool if __name__ == '__main__': wb_path = "你的目录路径" # 你的分块代码 entries = os.listdir(wb_path) df = pd.DataFrame(entries, columns=['Wk_nm']) n = 50 list_df = [df[i:i+n] for i in range(0, len(df), n)] # 初始化进程池(默认使用CPU核心数) with Pool() as pool: # 使用starmap传递多个参数(分块和目录路径) results = pool.starmap(process_chunk, [(chunk, wb_path) for chunk in list_df]) # 合并结果 all_results = pd.concat(results, ignore_index=True) print(all_results)
方法3:使用joblib(适合批量任务,支持进度条)
如果你需要查看处理进度,可以用joblib,需要先安装:pip install joblib
from joblib import Parallel, delayed if __name__ == '__main__': wb_path = "你的目录路径" # 你的分块代码 entries = os.listdir(wb_path) df = pd.DataFrame(entries, columns=['Wk_nm']) n = 50 list_df = [df[i:i+n] for i in range(0, len(df), n)] # 并行处理,n_jobs指定进程数,verbose显示进度 results = Parallel(n_jobs=-1, verbose=10)( delayed(process_chunk)(chunk, wb_path) for chunk in list_df ) all_results = pd.concat(results, ignore_index=True) print(all_results)
注意事项
- 必须把主逻辑放在
if __name__ == '__main__':下,避免多进程启动时重复执行代码 - 处理函数
process_chunk要能被序列化(pickle),所以不要在函数内部定义嵌套函数或使用无法序列化的对象 - 根据你的CPU核心数调整进程数,默认会使用全部核心,也可以手动指定(比如
ProcessPoolExecutor(max_workers=4))
内容的提问来源于stack exchange,提问作者Suman
相关产品推荐
相关产品推荐

