macOS下获取锁时间歇性出现ConnectionRefusedError问题求助
多进程共享内存锁引发ConnectionRefusedError问题分析与解决
问题现象
使用多进程操作共享内存时,当任务数RANGE较低(如25)程序稳定运行,但增至100左右时,间歇性在获取锁时抛出ConnectionRefusedError(Errno 61),最终导致计数结果不符合预期。
复现代码
from concurrent.futures import ProcessPoolExecutor from multiprocessing.shared_memory import SharedMemory from multiprocessing import Manager from struct import unpack_from, pack_into from itertools import repeat FMT = '@I' RANGE = 100 SHM = '_shm' class ShmWm(SharedMemory): def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) def __enter__(self): return self def __exit__(self, *_): super().close() def unlink(self): super().unlink() def process(lock): try: locked = False with lock: locked = True with ShmWm(SHM) as shared: v, = unpack_from(FMT, shared.buf) pack_into(FMT, shared.buf, 0, v+1) except ConnectionRefusedError as e: if not locked: print('Failed to acquire lock', e) raise def main(): with Manager() as manager: lock = manager.Lock() try: with ShmWm(SHM, True, 4) as shared: try: pack_into(FMT, shared.buf, 0, 0) with ProcessPoolExecutor() as executor: executor.map(process, repeat(lock, RANGE)) v, = unpack_from(FMT, shared.buf) assert v == RANGE print('All OK') finally: shared.unlink() except AssertionError: print(f'{v=} but was expected to be {RANGE}') except Exception as e: print('main:', e) if __name__ == '__main__': main()
运行平台
Python 3.10.6, macOS 12.5.1, 3 GHz 10-Core Intel Xeon W, 32 GB 2666 MHz DDR4
错误输出
Failed to acquire lock [Errno 61] Connection refused Process SpawnProcess-15: Failed to acquire lock [Errno 61] Connection refused Traceback (most recent call last): File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/process.py", line 314, in _bootstrap self.run() File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/process.py", line 108, in run self._target(*self._args, **self._kwargs) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/concurrent/futures/process.py", line 240, in _process_worker call_item = call_queue.get(block=True) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/queues.py", line 122, in get return _ForkingPickler.loads(res) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 942, in RebuildProxy return func(token, serializer, incref=incref, **kwds) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 792, in __init__ self._incref() File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 846, in _incref conn = self._Client(self._token.address, authkey=self._authkey) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/connection.py", line 507, in Client c = SocketClient(address) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/connection.py", line 635, in SocketClient s.connect(address) ConnectionRefusedError: [Errno 61] Connection refused Failed to acquire lock [Errno 61] Connection refused Process SpawnProcess-2: Traceback (most recent call last): File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/process.py", line 314, in _bootstrap self.run() File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/process.py", line 108, in run self._target(*self._args, **self._kwargs) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/concurrent/futures/process.py", line 240, in _process_worker call_item = call_queue.get(block=True) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/queues.py", line 122, in get return _ForkingPickler.loads(res) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 942, in RebuildProxy return func(token, serializer, incref=incref, **kwds) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 792, in __init__ self._incref() File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/managers.py", line 846, in _incref conn = self._Client(self._token.address, authkey=self._authkey) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/connection.py", line 507, in Client c = SocketClient(address) File "/Library/Frameworks/Python.framework/Versions/3.10/lib/python3.10/multiprocessing/connection.py", line 635, in SocketClient s.connect(address) ConnectionRefusedError: [Errno 61] Connection refused v=65 but was expected to be 100
问题原因
- Manager锁的通信机制缺陷:
Manager.Lock是基于Socket通信的代理对象,所有子进程需要连接Manager服务进程来获取锁。当并发进程数过多时,Manager服务的Socket连接池会被耗尽,或者服务端无法及时处理大量连接请求,导致连接被拒绝。 - macOS的Spawn启动方式:macOS下Python多进程默认使用Spawn方式,每个子进程都需要重新初始化并连接Manager服务,高并发下这种连接请求的竞争会加剧,容易出现连接失败。
- 进程池任务调度的影响:当任务数远大于进程池的默认进程数时,任务会排队等待,子进程重复使用时可能出现代理对象的连接失效问题。
解决思路与修改方案
方案1:使用原生multiprocessing.Lock替代Manager.Lock
multiprocessing.Lock是基于操作系统原生同步原语实现的,不需要通过Socket通信的Manager服务,直接在进程间共享,从根本上避免了连接拒绝问题。
修改后的代码:
from concurrent.futures import ProcessPoolExecutor from multiprocessing.shared_memory import SharedMemory from multiprocessing import Lock # 替换Manager为Lock from struct import unpack_from, pack_into from itertools import repeat FMT = '@I' RANGE = 100 SHM = '_shm' class ShmWm(SharedMemory): def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) def __enter__(self): return self def __exit__(self, *_): super().close() def unlink(self): super().unlink() def process(lock): try: locked = False with lock: locked = True with ShmWm(SHM) as shared: v, = unpack_from(FMT, shared.buf) pack_into(FMT, shared.buf, 0, v+1) except Exception as e: if not locked: print('Failed to acquire lock', e) raise def main(): lock = Lock() # 直接创建原生Lock try: with ShmWm(SHM, True, 4) as shared: try: pack_into(FMT, shared.buf, 0, 0) with ProcessPoolExecutor() as executor: executor.map(process, repeat(lock, RANGE)) v, = unpack_from(FMT, shared.buf) assert v == RANGE print('All OK') finally: shared.unlink() except AssertionError: print(f'{v=} but was expected to be {RANGE}') except Exception as e: print('main:', e) if __name__ == '__main__': main()
方案2:限制进程池并发数
如果必须使用Manager.Lock,可以通过设置ProcessPoolExecutor的max_workers参数,减少同时连接Manager服务的进程数,缓解服务端压力。例如:
with ProcessPoolExecutor(max_workers=10) as executor: # 根据CPU核心数调整 executor.map(process, repeat(lock, RANGE))
方案3:优化锁的使用逻辑
减少锁的持有时间,避免在锁内执行不必要的操作(比如原代码中打开共享内存的操作可以移到锁外,因为共享内存本身是进程安全的,只有读写数据时需要锁):
def process(lock): try: with ShmWm(SHM) as shared: with lock: v, = unpack_from(FMT, shared.buf) pack_into(FMT, shared.buf, 0, v+1) except ConnectionRefusedError as e: print('Failed to acquire lock', e) raise
内容的提问来源于stack exchange,提问作者jackal
相关产品推荐
相关产品推荐

