Python班级类别并行计算实现求助:串行代码转并行方案咨询
串行代码改造为并行计算方案
核心思路
原代码中每个班级的数据筛选、图表计算逻辑完全独立,不存在依赖关系,这是并行化的理想场景。我们可以借助Python的concurrent.futures库,用多进程(CPU密集型任务优先)或多线程(IO密集型任务优先)实现并行处理,大幅提升计算效率。
改造步骤与代码示例
1. 封装单班级处理逻辑
把单个班级的完整处理流程(从数据筛选到图表数据生成)封装成独立函数,消除全局变量依赖:
import numpy as np from datetime import datetime, date from concurrent.futures import ProcessPoolExecutor def process_single_class(classes, df_all, student_info_df, charts_list): classes = int(classes) # 筛选当前班级的基础数据 class_df = df_all[df_all['Class'] == classes] student_num = len(df_all['student_id'][df_all['Class'] == classes].unique()) student_info_temp = student_info_df[student_info_df['Class'] == classes] # 初始化图表数据容器 student_index = 0 X_values = [[] for _ in range(len(charts_list))] Y_values = [[] for _ in range(len(charts_list))] labels = [[] for _ in range(len(charts_list))] data_all = [[[]] for _ in range(len(charts_list))] # 处理各类图表数据 for chart, chart_index in zip(charts_list, range(len(charts_list))): if chart == "Future Admission": data_all[chart_index][student_index] = [[], [], []] X_values[chart_index], data_all[chart_index][student_index][0], data_all[chart_index][student_index][1], data_all[chart_index][student_index][2] = get_activation_data(class_df, date.today()) data_all[chart_index][student_index] = [[(np.array(data_all[chart_index][student_index][i]) * student_num).tolist() for i in range(len(data_all[chart_index][student_index]))], student_num] labels[chart_index] = ['Admission', 'Retention', 'Activation'] elif chart == "Future Revenue Projection": data_all[chart_index][student_index] = [[], []] temp_start_date = str(datetime.today().date()) X_values[chart_index], data_all[chart_index][student_index][0], data_all[chart_index][student_index][1] = get_revenue_projection(class_df, date.today()) labels[chart_index] = ['25th Percentile', '50th Percentile', '75th Percentile', 'Median Revenue', '90th Percentile'] else: print(f"未知图表类型: {chart}") # 返回当前班级的所有处理结果 return { 'class_id': classes, 'class_df': class_df, 'student_info_temp': student_info_temp, 'student_num': student_num, 'X_values': X_values, 'Y_values': Y_values, 'labels': labels, 'data_all': data_all }
2. 并行执行所有班级的处理
替换原有的串行循环,用进程池批量处理所有班级:
# 定义全局配置(也可作为参数传入) charts_list = ['Future Admission', 'Future Revenue Projection'] if __name__ == '__main__': # 多进程必须添加该判断,避免子进程重复初始化 # 初始化进程池,max_workers可设为CPU核心数,或根据实际需求调整 with ProcessPoolExecutor(max_workers=4) as executor: # 提交所有班级的处理任务 futures = [executor.submit(process_single_class, cls, df_all, student_info_df, charts_list) for cls in class_list] # 收集所有任务结果 results = [] for future in futures: try: results.append(future.result()) except Exception as e: print(f"处理班级时出错: {e}") # 将结果整理为原代码的class_groups格式 class_groups = {} for res in results: class_id = res.pop('class_id') class_groups[class_id] = res
关键注意事项
- 池类型选择:如果
get_activation_data和get_revenue_projection是CPU密集型计算(如大量数值运算),用ProcessPoolExecutor;如果是IO密集型(如调用外部API、读写文件),换用ThreadPoolExecutor更高效。 - 数据序列化:多进程中传递的数据需支持pickle序列化,Pandas DataFrame和NumPy数组默认支持,若有自定义对象需确保其可序列化。
- 异常处理:并行任务中单个班级出错不会中断全局流程,通过
future.result()捕获异常可定位问题班级。 - 内存控制:若班级数量多、单班级数据量大,可分批提交任务,避免一次性占用过多内存。
内容的提问来源于stack exchange,提问作者r ram
相关产品推荐
相关产品推荐

