如何在Python中对同一数据的多个子集并行执行相同逻辑
Python实现按门店编号并行处理数据集子集
当然可以通过多进程实现多个门店子集的并行处理,以下是两种常用的实现方案:
方案1:使用multiprocessing.Pool
这是Python标准库中最常用的多进程工具之一,适合CPU密集型的处理任务:
import pandas as pd from multiprocessing import Pool # 定义你的复杂处理函数,输入为(门店编号, 子集数据),输出处理结果 def process_store(store_tuple): store_num, subset = store_tuple # 这里替换成你的实际复杂逻辑 total_sales = subset['sales_amount'].sum() # 其他复杂计算、数据清洗或分析步骤... return (store_num, total_sales) if __name__ == '__main__': # 加载你的数据集 data1 = pd.read_csv("your_data_source.csv") # 按store_number分组,转换为可迭代的(门店编号, 子集)元组列表 store_groups = list(data1.groupby('store_number')) # 创建进程池,processes参数建议设为CPU核心数(可通过os.cpu_count()获取) with Pool(processes=4) as pool: # 并行映射处理每个门店子集 processed_results = pool.map(process_store, store_groups) # 将结果整理成字典或DataFrame,方便后续使用 result_dict = dict(processed_results) print(result_dict)
方案2:使用concurrent.futures.ProcessPoolExecutor
这个API更简洁,语法更直观,同样基于多进程实现:
import pandas as pd from concurrent.futures import ProcessPoolExecutor def process_store(store_tuple): store_num, subset = store_tuple # 替换为你的实际复杂处理逻辑 total_sales = subset['sales_amount'].sum() # 其他操作... return (store_num, total_sales) if __name__ == '__main__': data1 = pd.read_csv("your_data_source.csv") store_groups = list(data1.groupby('store_number')) # 创建进程执行器,max_workers设置并发进程数 with ProcessPoolExecutor(max_workers=4) as executor: processed_results = list(executor.map(process_store, store_groups)) # 转换为DataFrame格式 result_df = pd.DataFrame(processed_results, columns=['store_number', 'total_sales']) print(result_df)
关键注意事项
- 进程数设置:
max_workers或processes建议设为CPU核心数(可通过os.cpu_count()获取),过多的进程会导致上下文切换开销剧增,反而降低效率。 - 函数序列化:处理函数和传递的数据必须是可被
pickle序列化的,否则多进程间无法传递数据会报错。如果函数依赖外部资源,建议在函数内部初始化。 - 大数据优化:如果数据集极大,分组后的数据传递可能占用大量内存,可考虑先将每个门店的子集保存为独立文件,让每个进程读取对应文件处理,减少内存压力。
内容的提问来源于stack exchange,提问作者Kenneth Singh
相关产品推荐
相关产品推荐

