【问题标题】:Why does random.sample used with multiprocessing.Pool deadlock sometimes?为什么 random.sample 有时与 multiprocessing.Pool 一起使用死锁?
【发布时间】:2020-02-26 20:21:21
【问题描述】:

当我运行以下 sn-p 时,有时它会死锁并且不会完成,但有时会。为什么会这样?我在 Ubuntu 16.04(4.4.0-173-generic)上运行 Python 3.8。

from functools import partial
from multiprocessing.pool import Pool
from random import sample

pool = Pool(4)
result = pool.map(partial(sample, range(10)), range(10))

当我为每个函数调用创建一个新的 random.Random 实例时,也会发生同样的情况:

def sample(data, k):
    rand = random.Random()
    return rand.sample(data, k)

如果它挂起并且我发送 SIGINT 我会得到以下回溯,但是我无法理解它:

Error in atexit._run_exitfuncs:
Traceback (most recent call last):
  File "[...]/multiprocessing/util.py", line 277, in _run_finalizers
Process ForkPoolWorker-2:
Process ForkPoolWorker-1:
Process ForkPoolWorker-3:
    finalizer()
  File "[...]/multiprocessing/util.py", line 201, in __call__
    res = self._callback(*self._args, **self._kwargs)
  File "[...]/multiprocessing/pool.py", line 689, in _terminate_pool
    cls._help_stuff_finish(inqueue, task_handler, len(pool))
  File "[...]/multiprocessing/pool.py", line 674, in _help_stuff_finish
    inqueue._rlock.acquire()
KeyboardInterrupt
Traceback (most recent call last):
Traceback (most recent call last):
  File "[...]/multiprocessing/process.py", line 315, in _bootstrap
    self.run()
  File "[...]/multiprocessing/process.py", line 315, in _bootstrap
    self.run()
  File "[...]/multiprocessing/process.py", line 108, in run
    self._target(*self._args, **self._kwargs)
  File "[...]/multiprocessing/process.py", line 108, in run
    self._target(*self._args, **self._kwargs)
  File "[...]/multiprocessing/pool.py", line 114, in worker
    task = get()
  File "[...]/multiprocessing/pool.py", line 114, in worker
    task = get()
  File "[...]/multiprocessing/queues.py", line 355, in get
    with self._rlock:
  File "[...]/multiprocessing/queues.py", line 355, in get
    with self._rlock:
  File "[...]/multiprocessing/synchronize.py", line 95, in __enter__
    return self._semlock.__enter__()
  File "[...]/multiprocessing/synchronize.py", line 95, in __enter__
    return self._semlock.__enter__()
KeyboardInterrupt
KeyboardInterrupt
Traceback (most recent call last):
  File "[...]/multiprocessing/process.py", line 315, in _bootstrap
    self.run()
  File "[...]/multiprocessing/process.py", line 108, in run
    self._target(*self._args, **self._kwargs)
  File "[...]/multiprocessing/pool.py", line 114, in worker
    task = get()
  File "[...]/multiprocessing/queues.py", line 356, in get
    res = self._reader.recv_bytes()
  File "[...]/multiprocessing/connection.py", line 216, in recv_bytes
    buf = self._recv_bytes(maxlength)
  File "[...]/multiprocessing/connection.py", line 414, in _recv_bytes
    buf = self._recv(4)
  File "[...]/multiprocessing/connection.py", line 379, in _recv
    chunk = read(handle, remaining)
KeyboardInterrupt
Process ForkPoolWorker-4:
Traceback (most recent call last):
  File "[...]/multiprocessing/process.py", line 315, in _bootstrap
    self.run()
  File "[...]/multiprocessing/process.py", line 108, in run
    self._target(*self._args, **self._kwargs)
  File "[...]/multiprocessing/pool.py", line 114, in worker
    task = get()
  File "[...]/multiprocessing/queues.py", line 355, in get
    with self._rlock:
  File "[...]/multiprocessing/synchronize.py", line 95, in __enter__
    return self._semlock.__enter__()
KeyboardInterrupt

【问题讨论】:

  • 共享资源需要被保护,以免被另一个线程(进程)访问时已经被另一个线程(进程)访问。如果您实现了线程锁,请确保在完成共享资源后解锁。
  • @ryyker 你能提供更多细节吗?共享了哪些资源,尤其是在每个函数调用创建一个新的 random.Random 实例时?
  • 嗯,这对我来说似乎是 Python 3.8 的错误。我可以重现 3.8.2 上的挂起,但不能重现 3.7.3。
  • 我没有python,或者我没有Linux,所以不能评论太多,但是inqueue._rlock.acquire()是否需要与伴随调用配对才能释放锁?当我看到lock 时,我会觉得它是为了保护某些东西,或者阻止对某些东西的访问。除非在其他进程需要访问某个东西时从锁中释放它,否则事情将开始堆积。如果这种情况发生的时间过长,可能会出现问题。
  • ...不过,@dano 可能是对的,有这个bug 报告。

标签: python python-3.x multiprocessing python-multiprocessing


【解决方案1】:

这是 Python 3.8 中的 open bug。它与random 无关,原因是工作进程的清理似乎没有正确完成。例如以下死锁:

from multiprocessing.pool import Pool

def test(x):
    return 'test'

pool = Pool(4)
result = pool.map(test, range(10))

解决方案是在 map 返回后手动调用 pool.close() 或使用池对象作为上下文管理器:

with Pool(4) as pool:
    result = pool.map(test, range(10))

【讨论】:

    猜你喜欢
    • 2013-01-14
    • 2021-03-14
    • 1970-01-01
    • 2011-04-18
    • 2011-08-31
    • 2022-11-26
    • 1970-01-01
    • 1970-01-01
    • 2020-11-01
    相关资源
    最近更新 更多