代码之家  ›  专栏  ›  技术社区  ›  s-m-e

尝试在多进程中访问持久数据时获取不稳定的运行时异常。池工作进程

  •  1
  • s-m-e  · 技术社区  · 7 年前

    受到启发 this solution 我正试图在Python中建立一个多进程工作进程池。其思想是在工作进程实际开始工作并最终重用之前将一些数据传递给它们。它的目的是将每次调用都需要打包/解包到工作进程中的数据量最小化(即减少进程间通信开销)。我的 MCVE 如下所示:

    import multiprocessing as mp
    import numpy as np
    
    def create_worker_context():
        global context # create "global" context in worker process
        context = {}
    
    def init_worker_context(worker_id, some_const_array, DIMS, DTYPE):
        context.update({
            'worker_id': worker_id,
            'some_const_array': some_const_array,
            'tmp': np.zeros((DIMS, DIMS), dtype = DTYPE),
            }) # store context information in global namespace of worker
        return True # return True, verifying that the worker process received its data
    
    class data_analysis:
        def __init__(self):
            self.DTYPE = 'float32'
            self.CPU_LEN = mp.cpu_count()
            self.DIMS = 100
            self.some_const_array = np.zeros((self.DIMS, self.DIMS), dtype = self.DTYPE)
            # Init multiprocessing pool
            self.cpu_pool = mp.Pool(processes = self.CPU_LEN, initializer = create_worker_context) # create pool and context in workers
            pool_results = [
                self.cpu_pool.apply_async(
                    init_worker_context,
                    args = (core_id, self.some_const_array, self.DIMS, self.DTYPE)
                ) for core_id in range(self.CPU_LEN)
                ] # pass information to workers' context
            result_batches = [result.get() for result in pool_results] # check if they got the information
            if not all(result_batches): # raise an error if things did not work
                raise SyntaxError('Workers could not be initialized ...')
    
        @staticmethod
        def process_batch(batch_data):
            context['tmp'][:,:] = context['some_const_array'] + batch_data # some fancy computation in worker
            return context['tmp'] # return result
    
        def process_all(self):
            input_data = np.arange(0, self.DIMS ** 2, dtype = self.DTYPE).reshape(self.DIMS, self.DIMS)
            pool_results = [
                self.cpu_pool.apply_async(
                    data_analysis.process_batch,
                    args = (input_data,)
                ) for _ in range(self.CPU_LEN)
                ] # let workers actually work
            result_batches = [result.get() for result in pool_results]
            for batch in result_batches[1:]:
                np.add(result_batches[0], batch, out = result_batches[0]) # reduce batches
            print(result_batches[0]) # show result
    
    if __name__ == '__main__':
        data_analysis().process_all()
    

    我正和CPython 3.6.6一起运行上述内容。

    奇怪的是…有时是有效的,有时不是。如果不起作用, process_batch 方法引发异常,因为它找不到 some_const_array 作为一个关键 context . 完整的回溯过程如下:

    (env) me@box:/path> python so.py 
    multiprocessing.pool.RemoteTraceback: 
    """
    Traceback (most recent call last):
      File "/python3.6/multiprocessing/pool.py", line 119, in worker
        result = (True, func(*args, **kwds))
      File "/path/so.py", line 37, in process_batch
        context['tmp'][:,:] = context['some_const_array'] + batch_data # some fancy computation in worker
    KeyError: 'some_const_array'
    """
    
    The above exception was the direct cause of the following exception:
    
    Traceback (most recent call last):
      File "/path/so.py", line 54, in <module>
        data_analysis().process_all()
      File "/path/so.py", line 48, in process_all
        result_batches = [result.get() for result in pool_results]
      File "/path/so.py", line 48, in <listcomp>
        result_batches = [result.get() for result in pool_results]
      File "/python3.6/multiprocessing/pool.py", line 644, in get
        raise self._value
    KeyError: 'some_const_array'
    

    我很困惑。这是怎么回事?

    如果我 语境 字典包含“更高类型”的对象,例如数据库驱动程序或类似的对象,我不会遇到这种问题。只有当我 上下文 字典包含基本的python数据类型、集合或numpy数组。

    (是否有一种潜在的更好的方法以更可靠的方式实现相同的目标?我知道我的方法被认为是 hack ……)

    1 回复  |  直到 7 年前
        1
  •  1
  •   Darkonaut    7 年前

    您需要重新定位 init_worker_context 进入你 initializer 功能 create_worker_context .

    你的假设是 每一个 工作进程将运行 初始化工作环境 对你的困惑负责。 提交到池的任务将被送入一个内部任务队列,所有工作进程都从中读取。在您的情况下,一些工作进程完成了它们的任务,并再次竞争以获得新任务。因此,一个工作进程将执行多个任务,而另一个工作进程将无法获得单个任务。记住操作系统为线程(工作进程)调度运行时。