You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 01:44:53