You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 11:53:15