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
相关产品推荐
相关产品推荐

