Python多进程池全局计数器实现报错求助
多进程统计已处理项目数量的RuntimeError解决方法
我想要统计已处理的项目数量,但运行代码时出现错误,以下是最小复现代码及错误堆栈信息:
import multiprocessing from multiprocessing import Value from ctypes import c_int class DemoProcessor(object): def __init__(self, data_to_process): self.counter = Value(c_int, 0) self.data_to_process = data_to_process def _process_data(self, item): print (f"Processing {item}") self.counter.value = self.counter.value + 1 print (f"Processed {self.counter.value}") def process(self): pool = multiprocessing.Pool(processes=4) results_multi = pool.starmap(self._process_data, self.data_to_process), if __name__ == '__main__': demo = DemoProcessor(data_to_process=list(range(10))) demo.process()
错误信息:
Traceback (most recent call last): File "demo.py", line 23, in <module> demo.process() File "demo.py", line 19, in process results_multi = pool.starmap(self._process_data, self.data_to_process), File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\pool.py", line 372, in starmap return self._map_async(func, iterable, starmapstar, chunksize).get() File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\pool.py", line 771, in get raise self._value File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\pool.py", line 537, in _handle_tasks put(task) File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\connection.py", line 211, in send self._send_bytes(_ForkingPickler.dumps(obj)) File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\reduction.py", line 51, in dumps cls(buf, protocol).dump(obj) File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\sharedctypes.py", line 199, in __reduce__ assert_spawning(self) File "C:\Users\user\AppData\Local\Programs\Python\Python39\lib\multiprocessing\context.py", line 359, in assert_spawning raise RuntimeError( RuntimeError: Synchronized objects should only be shared between processes through inheritance
错误提示说明:同步对象应仅通过继承在进程间共享。
错误原因
问题出在multiprocessing.Value的传递方式上:当使用Pool.starmap时,会把类实例的_process_data方法连同整个实例对象序列化后传递给子进程,但Value这类底层同步对象无法被序列化,只能通过「进程继承」的方式共享——也就是在子进程启动前就创建好,让子进程从父进程的内存空间中直接继承该对象。
修复方案
方案1:使用Pool的initializer传递共享计数器
通过Pool的initializer和initargs参数,在子进程启动时初始化全局共享变量,让子进程继承Value对象:
import multiprocessing from multiprocessing import Value from ctypes import c_int # 全局变量,用于子进程访问共享计数器 counter = None def init_counter(shared_counter): global counter counter = shared_counter def _process_data(item): print(f"Processing {item}") counter.value += 1 print(f"Processed {counter.value}") class DemoProcessor(object): def __init__(self, data_to_process): self.counter = Value(c_int, 0) self.data_to_process = data_to_process def process(self): # 创建进程池时初始化子进程的全局计数器 pool = multiprocessing.Pool(processes=4, initializer=init_counter, initargs=(self.counter,)) # starmap要求每个任务是可迭代对象,因此包装成元组 pool.starmap(_process_data, [(item,) for item in self.data_to_process]) pool.close() pool.join() print(f"总处理项目数: {self.counter.value}") if __name__ == '__main__': demo = DemoProcessor(data_to_process=list(range(10))) demo.process()
方案2:使用Manager创建可序列化的共享对象
multiprocessing.Manager创建的共享对象通过代理实现,可以安全地在进程间传递,无需依赖继承:
import multiprocessing def _process_data(item, counter): print(f"Processing {item}") counter.value += 1 print(f"Processed {counter.value}") class DemoProcessor(object): def __init__(self, data_to_process): # 使用Manager创建可传递的共享计数器 manager = multiprocessing.Manager() self.counter = manager.Value('i', 0) self.data_to_process = data_to_process def process(self): pool = multiprocessing.Pool(processes=4) # 将计数器作为任务参数传递给每个子进程 tasks = [(item, self.counter) for item in self.data_to_process] pool.starmap(_process_data, tasks) pool.close() pool.join() print(f"总处理项目数: {self.counter.value}") if __name__ == '__main__': demo = DemoProcessor(data_to_process=list(range(10))) demo.process()
核心总结
- 直接用
multiprocessing.Value创建的同步对象属于父进程内存,不能序列化传递给子进程; - 共享这类对象的两种合法方式:
- 子进程启动前创建对象,通过继承让子进程访问;
- 使用
multiprocessing.Manager创建代理对象,支持序列化传递。
内容的提问来源于stack exchange,提问作者Baron Yugovich
相关产品推荐
相关产品推荐

