【问题标题】:Python multiprocessing synchronizationPython 多进程同步
【发布时间】:2014-10-11 16:58:26
【问题描述】:

我有一个函数“函数”,我想调用 10 次,使用 2 次 5 cpus 和多处理。

因此,我需要一种方法来同步以下代码中描述的进程。

如果不使用多处理池,这可能吗?如果这样做,我会收到奇怪的错误(例如“UnboundLocalError:分配前引用的局部变量'fd'”(我没有这样的变量))。进程似乎也随机终止。

如果可能的话,我想在没有游泳池的情况下这样做。谢谢!

number_of_cpus = 5
number_of_iterations = 2

# An array for the processes.
processing_jobs = []

# Start 5 processes 2 times.
for iteration in range(0, number_of_iterations):

    # TODO SYNCHRONIZE HERE

    # Start 5 processes at a time.
    for cpu_number in range(0, number_of_cpus):

        # Calculate an offset for the current function call.
        file_offset = iteration * cpu_number * number_of_files_per_process

        p = multiprocessing.Process(target=function, args=(file_offset,))
        processing_jobs.append(p)
        p.start()

    # TODO SYNCHRONIZE HERE

这是我在池中运行代码时遇到的错误的(匿名)回溯:

Process Process-5:
Traceback (most recent call last):
  File "/usr/lib/python2.7/multiprocessing/process.py", line 258, in _bootstrap
    self.run()
  File "/usr/lib/python2.7/multiprocessing/process.py", line 114, in run
    self._target(*self._args, **self._kwargs)
  File "python_code_3.py", line 88, in function_x
    xyz = python_code_1.function_y(args)
  File "/python_code_1.py", line 254, in __init__
    self.WK =  file.WK(filename)
  File "/python_code_2.py", line 1754, in __init__
    self.__parse__(name, data, fast_load)
  File "/python_code_2.py", line 1810, in __parse__
    fd.close()
UnboundLocalError: local variable 'fd' referenced before assignment

大多数进程都这样崩溃,但不是全部。当我增加进程数量时,它们中的更多似乎崩溃了。我还认为这可能是由于内存限制...

【问题讨论】:

  • 你为什么不想使用Pool
  • 这个任务真的最适合Pool,如果你能让它正常工作,你愿意用一个吗?
  • 如果您提供完整的回溯,将会很有帮助。代码在 Linux 上运行良好。你用的是什么平台?
  • 另外,“同步”是什么意思?您的意思是在开始第二组之前等待第一组 5 个进程完成吗?是否需要将 function 中的任何内容返回给父级?
  • 将其编辑到您的问题中。

标签: python parallel-processing synchronization multiprocessing threadpool


【解决方案1】:

Pool 非常容易使用。这是一个完整的例子:

来源

import multiprocessing

def calc(num):
    return num*2

if __name__=='__main__':  # required for Windows
    pool = multiprocessing.Pool()   # one Process per CPU
    for output in pool.map(calc, [1,2,3]):
        print 'output:',output

输出

output: 2
output: 4
output: 6

【讨论】:

  • 您可能应该使用if __name__ == "__main__": 保护,这样它就可以在 Windows 上运行。
  • @shavenwarthog 是在正确的轨道上。你,@user2177047,需要想出一个更好的方法来共享在function 中打开的共同资源。最好的方法是在您的父线程中打开资源,将工作集中到您的衍生进程并报告给父进程(然后写入文件)。或者,这篇 SO 帖子可能会帮助您提供一些想法:stackoverflow.com/questions/659865/…
【解决方案2】:

以下是您可以在不使用池的情况下执行所需同步的方法:

import multiprocessing

def function(arg):
    print ("got arg %s" % arg)

if __name__ == "__main__":
    number_of_cpus = 5
    number_of_iterations = 2

    # An array for the processes.
    processing_jobs = []

    # Start 5 processes 2 times.
    for iteration in range(1, number_of_iterations+1):  # Start the range from 1 so we don't multiply by zero.

        # Start 5 processes at a time.
        for cpu_number in range(1, number_of_cpus+1):

            # Calculate an offset for the current function call.
            file_offset = iteration * cpu_number * number_of_files_per_process

            p = multiprocessing.Process(target=function, args=(file_offset,))
            processing_jobs.append(p)
            p.start()

        # Wait for all processes to finish.
        for proc in processing_jobs:
            proc.join()

        # Empty active job list.
        del processing_jobs[:]

        # Write file here
        print("Writing")

这里是Pool

import multiprocessing

def function(arg):
    print ("got arg %s" % arg)

if __name__ == "__main__":
    number_of_cpus = 5
    number_of_iterations = 2

    pool = multiprocessing.Pool(number_of_cpus)
    for i in range(1, number_of_iterations+1): # Start the range from 1 so we don't multiply by zero
        file_offsets = [number_of_files_per_process * i * cpu_num for cpu_num in range(1, number_of_cpus+1)] 
        pool.map(function, file_offsets)
        print("Writing")
        # Write file here

如您所见,Pool 解决方案更好。

不过,这并不能解决您的回溯问题。在不了解实际原因的情况下,我很难说如何解决这个问题。您可能需要使用multiprocessing.Lock 来同步对资源的访问。

【讨论】:

  • 谢谢。我将同时查看回溯和您的解决方案。一个问题。在外部 for 循环中删除 processing_jobs 后,我不必重新实例化它吗?
  • @user2177047 我们只是删除了processing_jobs 列表的所有内容,而不是对象本身。因此,您可以在下一次迭代时再次开始附加它。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-11-09
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-11-22
  • 1970-01-01
相关资源
最近更新 更多