如何用Python multiprocessing.pool实现带独立对象的并行Worker?
用multiprocessing.Pool实现进程绑定独立对象并执行方法
需求与问题
需要在Python中创建多个Worker进程,每个进程初始化一个独立对象;主进程向Worker发送命令及参数,使其执行自身对象的方法并并行返回结果。
尝试使用multiprocessing.pool实现时,即使设置进程数为N、chunksize=1,进程与参数的绑定仍呈随机状态。测试代码及输出如下:
测试代码
import multiprocessing def init(a): global myA myA = a def get_value(_): global myA return myA.value class A(): def __init__(self, value): self.value = value if __name__ == '__main__': N = 4 a_lst = [A(i) for i in range(N)] pool = multiprocessing.Pool(N) pool.map(init, a_lst, chunksize=1) print(pool.map(get_value, range(N), chunksize=1))
输出结果
[3, 1, 3, 1]
问题:能否用multiprocessing.pool实现该需求?若可以该怎么做?
问题原因分析
原代码的错误在于:
- 使用
pool.map(init, a_lst)分发初始化任务,虽然chunksize=1且进程数等于任务数,每个进程会处理一次初始化,但Pool的任务分配是动态随机的,第二次调用pool.map(get_value, ...)时,任务会被随机分配给已存在的进程。如果某个进程被重复分配任务,就会导致重复返回同一对象的结果。 - 这种方式无法保证每个进程仅绑定一个独立对象,也无法确保后续任务与初始化的对象一一对应。
解决方案:利用Pool的初始化机制
可以通过multiprocessing.Pool的initializer和initargs参数,在进程启动时完成对象的初始化,确保每个进程仅绑定一个独立对象。具体步骤如下:
- 创建一个队列,将所有需要初始化的独立对象放入队列。
- 定义进程初始化函数,每个进程启动时从队列中取出一个唯一对象,作为自身的全局实例。
- 定义执行方法的函数,通过方法名调用当前进程绑定对象的对应方法,并返回结果。
修正后的代码
import multiprocessing from queue import Empty def init_worker(queue): """进程初始化函数:从队列中获取唯一对象绑定到当前进程""" global myA try: myA = queue.get(block=False) except Empty: myA = None def execute_method(method_info): """执行当前进程绑定对象的指定方法""" global myA if not myA: return None method_name, args = method_info method = getattr(myA, method_name) return method(*args) class A(): def __init__(self, value): self.value = value def get_value(self): """示例方法:返回对象的value属性""" return self.value def add(self, num): """示例方法:对value执行加法操作""" return self.value + num if __name__ == '__main__': N = 4 # 创建N个独立的A对象 a_lst = [A(i) for i in range(N)] # 用队列分发初始化对象,每个进程启动时取一个 init_queue = multiprocessing.Queue() for a in a_lst: init_queue.put(a) # 创建Pool,指定进程初始化函数及参数 pool = multiprocessing.Pool(N, initializer=init_worker, initargs=(init_queue,)) # 测试:调用所有进程的get_value方法 task_list = [('get_value', ()) for _ in range(N)] results = pool.map(execute_method, task_list) print("get_value结果:", results) # 输出包含0-3的所有值,顺序可能因进程调度略有不同 # 测试:调用所有进程的add方法,传入不同参数 task_list = [('add', (10,)), ('add', (20,)), ('add', (30,)), ('add', (40,))] results = pool.map(execute_method, task_list) print("add方法结果:", results) # 输出对应计算结果,顺序可能略有不同 pool.close() pool.join()
说明
initializer=init_worker指定了进程启动时执行的初始化函数,initargs=(init_queue,)传递了初始化所需的队列参数。- 每个进程启动时会从队列中取出一个唯一的
A对象,绑定到自身的全局变量myA,后续所有任务都会复用该对象。 execute_method函数接收方法名和参数,通过getattr调用对象的对应方法,确保每个任务都使用当前进程绑定的独立对象执行逻辑。
内容的提问来源于stack exchange,提问作者Mads M Pedersen
相关产品推荐
相关产品推荐

