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

重新接入tmux会话时Python多进程脚本触发BrokenPipeError问题

问题描述

我在大学实验室服务器上运行长时间的Python脚本,为避免SSH连接中断导致进程终止,操作流程如下:

  • 启动tmux会话:
tmux new -s preprocessing
  • 在会话内运行脚本:
python3 path/to/my/scirpt.py
  • 退出会话并关闭SSH连接,一段时间后重新连接并接入会话:
tmux attach -t preprocessing

刚接入时能看到脚本的进度条输出:

Processing 1st Step:   7%|███▌                                                 | 61725/921600 [11:11:49<155:59:00,  1.53it/s]

但几秒后脚本抛出BrokenPipeError错误:

multiprocessing.pool.RemoteTraceback:
"""
Traceback (most recent call last):
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 125, in worker
    result = (True, func(*args, **kwds))
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 51, in starmapstar
    return list(itertools.starmap(args[0], args[1]))
  File "/home/plodir2/network-predictive-analysis/scripts/data_preprocessing_MILC+LAMMPS+UR_lammps_1056.py", line 419, in process_chunk
    if (application_iteration, router_id) not in list(processed_routers):
  File "<string>", line 2, in __len__
  File "/usr/lib/python3.10/multiprocessing/managers.py", line 817, in _callmethod
    conn.send((self._id, methodname, args, kwds))
  File "/usr/lib/python3.10/multiprocessing/connection.py", line 206, in send
    self._send_bytes(_ForkingPickler.dumps(obj))
  File "/usr/lib/python3.10/multiprocessing/connection.py", line 411, in _send_bytes
    self._send(header + buf)
  File "/usr/lib/python3.10/multiprocessing/connection.py", line 368, in _send
    n = write(self._handle, buf)
BrokenPipeError: [Errno 32] Broken pipe
"""

The above exception was the direct cause of the following exception:

Traceback (most recent call last):
  File "/home/plodir2/network-predictive-analysis/scripts/data_preprocessing_MILC+LAMMPS+UR_lammps_1056.py", line 600, in <module>
    pool.starmap(process_chunk, [(application_chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations) for application_chunk in application_chunks])
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 375, in starmap
    return self._map_async(func, iterable, starmapstar, chunksize).get()
  File "/usr/lib/python3.10/multiprocessing/pool.py", line 774, in get
    raise self._value
BrokenPipeError: [Errno 32] Broken pipe

关键观察

  • 报错仅在重新接入tmux会话后发生,进度条显示的11小时正好是脚本启动到我检查进度的时长。
  • 用小数据集测试时,保持SSH连接活跃未出现任何错误。

脚本核心逻辑

脚本将pandas的application_data数据集拆分为多块,通过400个进程并行处理以提升性能,使用multiprocessing.Manager创建进程间共享的数据结构:

...
if __name__ == "__main__":
    num_processes = 400
    application_chunks = np.array_split(application_data, num_processes)

    with Manager() as manager:
        data_to_add_list = manager.list()
        processed_routers = manager.list()
        inactive_nodes_lists_processed_iterations = manager.list()
        
        progress_counter = Value('i', 0)
        lock = Lock()

        with tqdm(total=len(application_data), desc="Processing 1st Step") as pbar:
            progress_thread = Thread(target=update_progress_bar, args=(len(application_data), pbar))
            progress_thread.start()

            with Pool(num_processes, initializer=init, initargs=(lock, progress_counter)) as pool:
                pool.starmap(process_chunk, [(application_chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations) for application_chunk in application_chunks])

            progress_thread.join()
...

其中processed_routers用于避免多个进程重复处理同一(application_iteration, router_id)元组,读写逻辑如下:

def process_chunk(chunk, data_to_add_list, processed_routers, inactive_nodes_lists_processed_iterations):

    for application_data_record in chunk.itertuples():
        application_iteration = application_data_record.iteration # 整数类型
        router_id = int(application_data_record.node // 4)
        ...

        add_router_data_ports = False
        with lock:
            if (application_iteration, router_id) not in list(processed_routers):
                # 标记该迭代下的路由器端口数据已处理
                processed_routers.append((application_iteration, router_id))
                add_router_data_ports = True

        if add_router_data_ports:
            # 执行具体处理逻辑
...
解决方案

1. 优化进程间共享数据结构的性能

当前每次判断(application_iteration, router_id)是否已处理时,都会将共享的processed_routers列表转为本地列表,这会触发跨进程的全量数据同步,当列表规模变大后,极易引发连接超时或断开。

改用共享字典存储已处理元组,判断操作的时间复杂度从O(n)降为O(1),大幅减少跨进程通信量:

# 替换Manager.list()为Manager.dict()
processed_routers = manager.dict()

# 修改process_chunk中的判断逻辑
with lock:
    key = (application_iteration, router_id)
    if key not in processed_routers:
        processed_routers[key] = True  # 用字典键存储元组,值可设为任意标识
        add_router_data_ports = True

2. 脱离终端依赖避免tmux重连干扰

tmux会话重连时会重置终端的标准输入输出关联,而tqdm进度条依赖终端文件描述符,这可能导致进程间通信的管道异常。可以通过两种方式解决:

  • 将输出重定向到日志文件:运行脚本时脱离终端,直接将输出写入文件
python3 path/to/my/scirpt.py > preprocessing.log 2>&1 &

之后查看进度只需读取日志:

tail -f preprocessing.log
  • 禁用tqdm的终端交互模式:在脚本中强制tqdm输出到标准输出,不依赖终端特性
with tqdm(total=len(application_data), desc="Processing 1st Step", disable=None, leave=True) as pbar:
    # 保持原进度更新逻辑

3. 调整进程数量至合理范围

400个进程远超多数实验室服务器的CPU核心数(通常为32-64核),会导致严重的进程调度开销,加剧跨进程通信的不稳定性。建议将进程数设置为CPU核心数的1-2倍:

import multiprocessing
num_processes = multiprocessing.cpu_count() * 2  # 例如32核服务器设置为64

4. 增强跨进程连接的稳定性

对于Python 3.8及以上版本,可以通过自定义Manager启用心跳机制,避免长时间无通信导致的连接断开:

from multiprocessing.managers import SyncManager

# 自定义Manager并注册所需数据类型
class CustomManager(SyncManager):
    pass

CustomManager.register('dict')
CustomManager.register('list')

# 初始化Manager时设置心跳间隔(单位:秒)
with CustomManager(address=('', 50000), authkey=b'secret_key', heartbeat_interval=30) as manager:
    # 原逻辑:创建共享数据结构、启动进程池等

内容的提问来源于stack exchange,提问作者Pietro Lodi Rizzini

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 04:47:32