Python多进程Pool.join()偶发[Errno 111]连接拒绝错误求助
Python multiprocessing map_async 偶发 ConnectionRefusedError 问题分析与解决
问题核心原因
从调用栈和代码来看,错误的根源是Manager进程与子进程/进程池的生命周期不匹配:
- 你在进程池外部创建了
Manager实例,其生命周期不受进程池约束。当进程池执行到join()阶段进行资源清理时,Manager进程可能已因为GC回收、意外退出或系统资源不足等原因提前终止。 - 进程池清理逻辑会尝试读取任务队列的残留数据,反序列化时需要重建
shared_dict的代理对象,但此时Manager进程已不可达,直接触发[Errno 111] Connection refused错误。
解决方案
方案1:用上下文管理器对齐Manager与进程池的生命周期
将Manager的创建放在with语句中,确保其生命周期完全覆盖进程池的运行、清理阶段:
with Manager() as manager: shared_dict = manager.dict() total_amount_of_data = 0 transfer_list_parameters = [] for item in transfer_list: total_amount_of_data += item[2] transfer_list_parameters.append((item, shared_dict)) total_data_mb = round(total_amount_of_data / (1024 * 1024), 2) with Pool(processes=copy_workers) as pool: result = pool.map_async(copyfile, transfer_list_parameters, chunksize=chunk_size, error_callback=copy_error_callback) while not result.ready(): print_progress() sleep(10) pool.close() pool.join()
方案2:替换共享字典为队列传递结果
避免直接传递Manager代理对象到子进程,改用Queue收集复制结果,减少跨进程连接的风险:
def copyfile(data_tuple): copy_info, result_queue = data_tuple src, dst, size, file_id = copy_info copy_op(src, dst) # 将结果放入队列,而非写入共享字典 result_queue.put((file_id, dict(size=size, src=src, dst=dst))) with Pool(processes=copy_workers) as pool: result_queue = Manager().Queue() transfer_list_parameters = [(item, result_queue) for item in transfer_list] result = pool.map_async(copyfile, transfer_list_parameters, chunksize=chunk_size, error_callback=copy_error_callback) while not result.ready(): print_progress() sleep(10) # 从队列收集所有结果到本地字典 shared_dict = {} while not result_queue.empty(): file_id, data = result_queue.get() shared_dict[file_id] = data pool.close() pool.join()
方案3:手动控制Manager的关闭时机
如果不想用上下文管理器,可通过try-finally确保进程池完全结束后再关闭Manager:
manager = Manager() try: shared_dict = manager.dict() total_amount_of_data = 0 transfer_list_parameters = [] for item in transfer_list: total_amount_of_data += item[2] transfer_list_parameters.append((item, shared_dict)) total_data_mb = round(total_amount_of_data / (1024 * 1024), 2) with Pool(processes=copy_workers) as pool: result = pool.map_async(copyfile, transfer_list_parameters, chunksize=chunk_size, error_callback=copy_error_callback) while not result.ready(): print_progress() sleep(10) pool.close() pool.join() finally: # 确保进程池完全结束后再关闭Manager manager.shutdown()
额外注意事项
- Python 3.7的multiprocessing模块在Manager代理的生命周期管理上存在已知缺陷,升级到Python 3.8及以上版本可能自动缓解该问题。
- 检查
copy_error_callback是否会抛出未处理异常,异常可能导致进程池提前终止,间接引发Manager连接问题。 - 监控系统资源使用情况,Manager进程可能因内存不足被系统OOM killer终止,从而触发连接失败。
内容的提问来源于stack exchange,提问作者Ismael Raya Roa
相关产品推荐
相关产品推荐

