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

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}

问题根源

  1. 重复初始化类:Windows系统下,multiprocess启动子进程时会重新导入整个脚本,脚本顶层的parameters = InputConect()会被每个子进程执行一次,因此出现多次输入提示。
  2. 子进程内存隔离:make_statistic试图修改主进程的字典,但子进程拥有独立内存空间,修改的是自身副本,主进程字典完全未更新,最终输出初始值0。
  3. 局部函数序列化隐患:嵌套的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 08:40:23