【问题标题】:How to use Python multiprocessing queue to access GPU (through PyOpenCL)?如何使用 Python 多处理队列访问 GPU(通过 PyOpenCL)?
【发布时间】:2015-06-19 04:52:56
【问题描述】:

我的代码需要很长时间才能运行,因此我一直在研究 Python 的多处理库以加快速度。我的代码还有几个步骤通过 PyOpenCL 利用 GPU。问题是,如果我将多个进程设置为同时运行,它们最终都会尝试同时使用 GPU,这通常会导致一个或多个进程抛出异常并退出。

为了解决这个问题,我错开每个进程的开始,这样它们就不太可能相互碰撞:

process_list = []
num_procs = 4

# break data into chunks so each process gets it's own chunk of the data
data_chunks = chunks(data,num_procs)
for chunk in data_chunks:
    if len(chunk) == 0:
        continue
    # Instantiates the process
    p = multiprocessing.Process(target=test, args=(arg1,arg2))
    # Sticks the thread in a list so that it remains accessible
    process_list.append(p)

# Start threads
j = 1
for process in process_list:
    print('\nStarting process %i' % j)
    process.start()
    time.sleep(5)
    j += 1

for process in process_list:
    process.join()

我还在调用 GPU 的函数周围包裹了一个 try except 循环,这样如果两个进程确实尝试同时访问它,则无法访问的进程将等待几秒钟并重试:

wait = 2
n = 0
while True:
    try:
        gpu_out = GPU_Obj.GPU_fn(params)
    except:
        time.sleep(wait)
        print('\n Waiting for GPU memory...')
        n += 1
        if n == 5:
            raise Exception('Tried and failed %i times to allocate memory for opencl kernel.' % n)
        continue
    break

这种解决方法非常笨拙,尽管它在大多数情况下都有效,但进程偶尔会抛出异常,我觉得应该有一个更有效/优雅的解决方案,使用 multiprocessing.queue 或类似的东西。但是,我不确定如何将它与 PyOpenCL 集成以进行 GPU 访问。

【问题讨论】:

    标签: python queue multiprocessing pyopencl


    【解决方案1】:

    听起来您可以使用multiprocessing.Lock 来同步对 GPU 的访问:

    data_chunks = chunks(data,num_procs)
    lock = multiprocessing.Lock()
    for chunk in data_chunks:
        if len(chunk) == 0:
            continue
        # Instantiates the process
        p = multiprocessing.Process(target=test, args=(arg1,arg2, lock))
        ...
    

    然后,在您访问 GPU 的 test 内部:

    with lock:  # Only one process will be allowed in this block at a time.
        gpu_out = GPU_Obj.GPU_fn(params)
    

    编辑:

    要对池执行此操作,您可以这样做:

    # At global scope
    lock = None
    
    def init(_lock):
        global lock
        lock = _lock
    
    data_chunks = chunks(data,num_procs)
    lock = multiprocessing.Lock()
    for chunk in data_chunks:
        if len(chunk) == 0:
            continue
        # Instantiates the process
        p = multiprocessing.Pool(initializer=init, initargs=(lock,))
        p.apply(test, args=(arg1, arg2))
        ...
    

    或者:

    data_chunks = chunks(data,num_procs)
    m = multiprocessing.Manager()
    lock = m.Lock()
    for chunk in data_chunks:
        if len(chunk) == 0:
            continue
        # Instantiates the process
        p = multiprocessing.Pool()
        p.apply(test, args=(arg1, arg2, lock))
    

    【讨论】:

    • 哇!这比我想象的要简单得多!如果您使用multiprocessing.Pool 而不是multiprocessing.Process,我假设您会以同样的方式使用multiprocessing.Lock
    • @johnny_be 如果使用池,实际上会有点复杂,因为您通常不能腌制Lock,因此将其传递给pool.map/pool.apply 将不起作用。在这种情况下,您需要将锁传递给Pool 构造函数:Pool(initializer=init, initargs=(lock,),然后在init 函数内使锁成为全局锁,或者使用可以腌制的multiprocessing.Manager().lock。有关更完整的示例,请参阅我的编辑。
    • 如果您只是同时在多个 python 控制台中运行相同的进程,我认为没有办法实现某种锁定? (即打开命令窗口、运行python my_script.py、打开另一个命令窗口、在新窗口中运行python my_script.py等)
    • @johnny_be 有一种方法可以使用multiprocessing.managers.BaseManager。请参阅this answer,它显示了如何使用Barrier 而不是Lock。不过,同样的想法也适用;您的脚本的一个实例是管理服务器,其他的是客户端。不过,您需要确保服务器实例在所有客户端完成之前不会退出。
    • @johnny_be 是的,只需创建另一个锁实例。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-05-15
    • 2013-10-02
    • 1970-01-01
    • 2015-08-14
    • 2021-02-12
    • 1970-01-01
    相关资源
    最近更新 更多