Python multiprocess多进程编程:类重复初始化与数据丢失问题
问题描述
最初遭遇Can't pickle local object错误,改用multiprocess库替代multiprocessing后,出现两个关键问题:
InputConect类的初始化次数与处理器核心数一致,导致重复弹出输入提示- 最终统计结果字典全为初始值0,数据未被正确计算或保存
原代码
import pathlib import pandas as pd from datetime import datetime import csv import re import multiprocess as mp class InputConect: def __init__(self): self.file_name = input('Введите название файла: ') self.filter_param = input('Введите название профессии: ') @staticmethod def print_data(file_name, filter_param): salary_by_years = {year: 0 for year in unique_years} vacs_by_years = {year: 0 for year in unique_years} vac_salary_by_years = {year: 0 for year in unique_years} vac_counts_by_years = {year: 0 for year in unique_years} def make_statistic(file): #writes data to the dictionary like this: salary_by_years[year] = int(one_year_vacancies.salary.mean()) vacs_by_years[year] = one_year_vacancies.shape[0] if __name__ == '__main__': with mp.Pool() as p: mp.freeze_support() p.map(make_statistic, filelist) p.close() p.join() # m = mp.map(target=make_statistic, args=filelist) # m.start() # m.join() print('Динамика уровня зарплат по годам:', salary_by_years) print('Динамика количества вакансий по годам:', vacs_by_years) print('Динамика уровня зарплат по годам для выбранной профессии:', vac_salary_by_years) print('Динамика количества вакансий по годам для выбранной профессии:', vac_counts_by_years) parameters = InputConect() InputConect.print_data(parameters.file_name, parameters.filter_param)
运行输出
Введите название файла: vacancies_by_year.csv
Введите название профессии: Аналитик
Введите название файла: Введите название файла: Введите название файла: Введите название файла: vacancies_by_year.csv
Введите название профессии: Аналитик
Введите название профессии: Аналитик
Введите название профессии: Аналитик
Введите название профессии: Аналитик
Динамика уровня зарплат по годам: {2007: 0, 2008: 0, 2009: 0, 2010: 0, 2011: 0, 2012: 0, 2013: 0, 2014: 0, 2015: 0, 2016: 0, 2017: 0, 2018: 0, 2019: 0, 2020: 0, 2021: 0, 2022: 0}
Динамика количества вакансий по годам: {2007: 0, 2008: 0, 2009: 0, 2010: 0, 2011: 0, 2012: 0, 2013: 0, 2014: 0, 2015: 0, 2016: 0, 2017: 0, 2018: 0, 2019: 0, 2020: 0, 2021: 0, 2022: 0}
Динамика уровня зарплат по годам для выбранной профессии: {2007: 0, 2008: 0, 2009: 0, 2010: 0, 2011: 0, 2012: 0, 2013: 0, 2014: 0, 2015: 0, 2016: 0, 2017: 0, 2018: 0, 2019: 0, 2020: 0, 2021: 0, 2022: 0}
Динамика количества вакансий по годам для выбранной профессии: {2007: 0, 2008: 0, 2009: 0, 2010: 0, 2011: 0, 2012: 0, 2013: 0, 2014: 0, 2015: 0, 2016: 0, 2017: 0, 2018: 0, 2019: 0, 2020: 0, 2021: 0, 2022: 0}
问题根源
- 重复初始化类:Windows系统下,
multiprocess启动子进程时会重新导入整个脚本,脚本顶层的parameters = InputConect()会被每个子进程执行一次,因此出现多次输入提示。 - 子进程内存隔离:
make_statistic试图修改主进程的字典,但子进程拥有独立内存空间,修改的是自身副本,主进程字典完全未更新,最终输出初始值0。 - 局部函数序列化隐患:嵌套的
make_statistic作为局部函数,即便用multiprocess也存在序列化风险,且访问外部变量的方式在多进程场景下不可靠。
修复方案
修复后代码
import pathlib import pandas as pd import multiprocess as mp class InputConect: def __init__(self): self.file_name = input('Введите название файла: ') self.filter_param = input('Введите название профессии: ') def make_statistic(args): file_path, filter_param = args # 读取单年份文件数据 df = pd.read_csv(file_path) # 从文件名提取年份(需根据实际文件名格式调整) year = int(pathlib.Path(file_path).stem.split('_')[-1]) # 计算全局统计数据 total_salary = int(df['salary'].mean()) if not df.empty else 0 total_count = df.shape[0] # 计算目标职业统计数据 filtered_df = df[df['name'].str.contains(filter_param, case=False, na=False)] filtered_salary = int(filtered_df['salary'].mean()) if not filtered_df.empty else 0 filtered_count = filtered_df.shape[0] # 返回当前文件的统计结果 return (year, total_salary, total_count, filtered_salary, filtered_count) def print_data(file_name, filter_param): # 获取所有目标文件列表 file_list = list(pathlib.Path('.').glob(f'*{file_name}*.csv')) # 初始化统计字典 unique_years = range(2007, 2023) salary_by_years = {y:0 for y in unique_years} vacs_by_years = {y:0 for y in unique_years} vac_salary_by_years = {y:0 for y in unique_years} vac_counts_by_years = {y:0 for y in unique_years} # 多进程处理文件 with mp.Pool() as pool: results = pool.map(make_statistic, [(file, filter_param) for file in file_list]) # 汇总子进程结果到字典 for year, sal, cnt, filt_sal, filt_cnt in results: if year in salary_by_years: salary_by_years[year] = sal vacs_by_years[year] = cnt vac_salary_by_years[year] = filt_sal vac_counts_by_years[year] = filt_cnt # 输出结果 print('Динамика уровня зарплат по годам:', salary_by_years) print('Динамика количества вакансий по годам:', vacs_by_years) print('Динамика уровня зарплат по годам для выбранной профессии:', vac_salary_by_years) print('Динамика количества вакансий по годам для выбранной профессии:', vac_counts_by_years) if __name__ == '__main__': mp.freeze_support() parameters = InputConect() print_data(parameters.file_name, parameters.filter_param)
修复核心要点
- 隔离顶层执行代码:将类实例化、函数调用放入
if __name__ == '__main__':块,避免子进程重复执行输入逻辑。 - 子进程返回结果:让
make_statistic返回单文件的统计数据,主进程统一汇总,规避内存隔离导致的变量修改无效问题。 - 顶层函数替代嵌套函数:将
make_statistic改为顶层函数,消除局部对象序列化风险,明确参数传递路径。
内容的提问来源于stack exchange,提问作者MpirtGod

