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

Python 2.7多进程中如何并行化类实例化并共享类对象?

在Python 2.7中用multiprocessing并行化MyContainer实例化是可行的,但需要重构代码解决共享问题

首先得明确:你的核心问题是多进程环境下类级变量(instances、children、OUTPUT_HEADINGS)无法跨进程同步——Python多进程基于fork(Unix)或spawn(Windows),每个子进程会复制一份主进程的内存空间,子进程对类变量的修改不会同步到主进程,这就是你遇到共享问题的根源。

下面是具体的解决方案,针对Python 2.7的特性调整:

1. 重构思路

  • 移除类级别的共享变量,改用multiprocessing.Manager提供的跨进程共享数据结构(比如Manager.list)来存储实例数据和表头信息。
  • 将文件处理逻辑从类的__new__方法中剥离出来,封装成独立函数,方便多进程调用。
  • 避免直接共享MyContainer实例(Python 2.7中对象pickle可能有兼容性问题),转而共享实例的字典表示(通过to_dict())。

2. 修改后的代码示例

先重构MyContainer类

import os
import logging
import csv
from collections import defaultdict
import xml.etree.ElementTree as ElementTree

class MyContainer(object):
    def __init__(self, dat_file, children):
        self._name = os.path.basename(dat_file)
        self.attr_value_sum = defaultdict(list)
        # 直接使用传入的children,不再依赖类变量
        var1_elem = children[0].find("var1")
        var1 = var1_elem.text if var1_elem is not None else ""
        var2 = children[0].get("var2", "")
        self.cat_name = "{}.{}".format(var1, var2)
        
        # 这里保留你原有的XML数据处理逻辑
        # ... 你的数据处理代码 ...

    def to_dict(self):
        # 返回处理后的字典,用于后续CSV导出
        output_dict = {
            # 填充你的字段,比如:
            self.cat_name: sum(self.attr_value_sum.values()),
            # 其他字段...
        }
        return output_dict

编写并行处理函数与主逻辑

from multiprocessing import Pool, Manager

def process_single_file(dat_file, shared_instances, shared_headings):
    """单个文件的处理逻辑,供多进程调用"""
    try:
        tree = ElementTree.parse(dat_file)
        children = tree.findall("parent_element/child_element")
    except ElementTree.ParseError as err:
        logging.exception(err)
        children = []
    
    if not children:
        msg = ("{}: No \"parent_element/child_element\" element found".format(os.path.basename(dat_file)))
        logging.warning(msg)
        return False
    
    # 创建实例并生成字典
    container = MyContainer(dat_file, children)
    instance_dict = container.to_dict()
    
    # 将数据存入共享结构
    shared_instances.append(instance_dict)
    if container.cat_name not in shared_headings:
        shared_headings.append(container.cat_name)
    
    return True

def main():
    # 初始化共享数据结构
    manager = Manager()
    shared_instances = manager.list()
    shared_headings = manager.list()

    # 创建进程池,默认用CPU核心数
    pool = Pool()
    process_results = []
    
    # 提交所有文件的处理任务
    for idx, f in enumerate(FILE_LIST, 1):
        print "{}/{}: {} in progress...".format(idx, len(FILE_LIST), f)
        res = pool.apply_async(process_single_file, args=(f, shared_instances, shared_headings))
        process_results.append((f, res))
    
    # 等待所有进程完成
    pool.close()
    pool.join()

    # 打印最终处理状态
    for f, res in process_results:
        print "{}...{}".format(f, "DONE" if res.get() else "SKIPPED")
    
    # 导出CSV
    output_headings = list(shared_headings)
    with open(args.output, "w") as output_file:
        f_csv = csv.DictWriter(output_file, fieldnames=output_headings)
        f_csv.writeheader()
        for instance_dict in shared_instances:
            f_csv.writerow(instance_dict)

if __name__ == '__main__':
    # 这里保留你原有的FILE_LIST初始化逻辑
    FILE_LIST = []
    for d in args.dirs:
        FILE_LIST.extend(get_name_defined_files(dir_path=d, pattern=args.filename, recursive=args.recursive))
    main()

3. 关键说明

  • multiprocessing.Manager会创建一个独立的服务器进程,所有子进程通过代理访问共享列表,保证了数据的跨进程同步。
  • 把文件处理逻辑放到独立函数中,避免了类变量带来的共享问题,每个子进程处理自己的文件,结果统一存入共享结构。
  • 用实例的字典表示代替直接共享实例,规避了Python 2.7中对象pickle的潜在问题,同时也更适合CSV导出的需求。

内容的提问来源于stack exchange,提问作者Petr Krampl

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 06:52:15