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

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创建的同步对象属于父进程内存,不能序列化传递给子进程;
  • 共享这类对象的两种合法方式:
    1. 子进程启动前创建对象,通过继承让子进程访问;
    2. 使用multiprocessing.Manager创建代理对象,支持序列化传递。

内容的提问来源于stack exchange,提问作者Baron Yugovich

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 14:57:04